#!/usr/bin/env python3
"""Move today's Distru sales orders to DELIVERING, with guardrails.

    python3 agent.py --dry-run --fixture sample-orders.json   # no key, no network: shows every decision
    python3 agent.py --dry-run                                # live read, no writes; needs DISTRU_API_TOKEN
    python3 agent.py                                         # live; writes, one at a time, logged

What it does, in order, for one delivery day:
  1. GET /orders filtered to today's delivery_datetime and the statuses we are allowed to move FROM.
  2. For each order: re-read it (find-before-write), check the transition is on the allowlist,
     check the ledger (idempotent: an order we already moved today is skipped), run the pre-write
     check, then POST /orders {"id": ..., "status": "DELIVERING"}.
  3. Stop when MAX_WRITES_PER_RUN is reached or on the first unexpected error. Print a summary.

Everything a reviewer needs to read is in the CONSTANTS block below.

Verified against Distru's public API docs and backend, September 2026:
  GET  /public/v1/orders?delivery_datetime=<after>,<before>&statuses[]=READY_TO_SHIP   (inclusive ISO 8601 range)
  GET  /public/v1/orders/{id}                                                          (eventually consistent, ~1 s)
  POST /public/v1/orders  {"id": "...", "status": "DELIVERING"}   sparse update: omitted fields stay as they are
  Order statuses: PENDING PROCESSING READY_TO_SHIP DELIVERING DELIVERED COMPLETED CANCELED
  DELIVERING requires every line fulfilled, a customer set, and a compliance transfer when any
  item is package-tracked. Distru rejects the write with a 400 otherwise; we log it and move on.

CC0 1.0. Version 1.0.0 (2026-09-17).
"""
import argparse
import datetime as dt
import json
import os
import sys
import time
import urllib.error
import urllib.parse
import urllib.request

# ============================ CONSTANTS: the reviewable part ============================
BASE_URL = "https://app.distru.com/public/v1"
TOKEN_ENV = "DISTRU_API_TOKEN"

# The only status changes this script may make. Deny by default: anything not here is skipped.
ALLOWED_TRANSITIONS = {
    "READY_TO_SHIP": {"DELIVERING"},
}
TARGET_STATUS = "DELIVERING"

MAX_WRITES_PER_RUN = 25          # stopping condition. Raise it after a week of clean runs.
WRITE_PAUSE_SECONDS = 0.5        # one write at a time, with a breath between
TIMEZONE_OFFSET_HOURS = -7       # your warehouse's UTC offset (Pacific daylight time = -7). "Today" is local.
LEDGER_PATH = "status-agent-ledger.jsonl"   # one JSON line per decision; this is also the idempotency record

# Pre-write gate. "deterministic" checks fields only. "typesafe" adds one Noul question
# (needs TYPESAFE_API_KEY and `pip install typesafe-sdk`); see lesson 06-05.
GATE_MODE = os.environ.get("STATUS_AGENT_GATE", "deterministic")
GATE_ASK_ABOVE = 0.30            # Noul P(true) for "something is off" above this -> skip and ask a human
# =========================================================================================


def headers(token):
    return {"Authorization": f"Bearer {token}", "Content-Type": "application/json", "Accept": "application/json"}


def masked(token):
    return "<" + TOKEN_ENV + ">" if not token else token[:4] + "..." + token[-4:]


def today_range_utc(now=None):
    """Local calendar day -> inclusive UTC range for delivery_datetime."""
    tz = dt.timezone(dt.timedelta(hours=TIMEZONE_OFFSET_HOURS))
    now = now or dt.datetime.now(tz)
    start = now.astimezone(tz).replace(hour=0, minute=0, second=0, microsecond=0)
    end = start + dt.timedelta(days=1) - dt.timedelta(microseconds=1)
    fmt = lambda d: d.astimezone(dt.timezone.utc).strftime("%Y-%m-%dT%H:%M:%S.%fZ")
    return fmt(start), fmt(end), start.date().isoformat()


def list_url(after, before):
    q = [("delivery_datetime", f"{after},{before}")] + [("statuses[]", s) for s in sorted(ALLOWED_TRANSITIONS)]
    return BASE_URL + "/orders?" + urllib.parse.urlencode(q)


def http(method, url, token, body=None):
    data = json.dumps(body).encode() if body is not None else None
    req = urllib.request.Request(url, data=data, headers=headers(token), method=method)
    try:
        with urllib.request.urlopen(req, timeout=60) as resp:
            return resp.status, json.loads(resp.read().decode("utf-8"))
    except urllib.error.HTTPError as e:
        try:
            payload = json.loads(e.read().decode("utf-8", errors="replace"))
        except Exception:
            payload = {"errors": [{"message": "unreadable error body", "pointer": ["base"]}]}
        return e.code, payload


def error_text(payload):
    return "; ".join(f"{'/'.join(map(str, x.get('pointer', [])))}: {x.get('message')}" for x in payload.get("errors", []))


# ----------------------------------- the ledger -----------------------------------
def ledger_load():
    done = set()
    if os.path.exists(LEDGER_PATH):
        with open(LEDGER_PATH, encoding="utf-8") as f:
            for line in f:
                try:
                    rec = json.loads(line)
                except json.JSONDecodeError:
                    continue
                if rec.get("action") == "wrote":
                    done.add((rec["order_id"], rec["day"], rec["to"]))
    return done


def ledger_write(rec, dry_run):
    rec["at"] = dt.datetime.now(dt.timezone.utc).isoformat()
    rec["dry_run"] = dry_run
    line = json.dumps(rec, ensure_ascii=False)
    print(line)
    if not dry_run:
        with open(LEDGER_PATH, "a", encoding="utf-8") as f:
            f.write(line + "\n")


# ----------------------------------- the gate -----------------------------------
def gate_deterministic(order):
    """Field checks only. Returns (ok, reason)."""
    if not order.get("company"):
        return False, "no customer on the order; DELIVERING needs one"
    items = order.get("items") or []
    if not items:
        return False, "order has no line items"
    unfulfilled = [i for i in items if not (i.get("package_id") or i.get("batch_id") or i.get("package") or i.get("batch"))]
    if unfulfilled and any(i.get("product_id") for i in unfulfilled):
        # Best effort: an item with only a product_id is unfulfilled per the API docs.
        return False, f"{len(unfulfilled)} line(s) look unfulfilled; Distru would reject DELIVERING"
    return True, "fields look complete"


def gate_typesafe(order, day):
    """One Noul question before the write. Shape follows TypeSafe's SDK as used in lesson 06-05."""
    try:
        from typesafe_sdk import Noul, TypeSafeClient  # type: ignore
    except ImportError:
        return None, "typesafe_sdk not installed; falling back to deterministic gate"
    client = TypeSafeClient()
    state = {"today": day, "order": {k: order.get(k) for k in ("order_number", "status", "delivery_datetime", "company", "total", "items")}}
    q = {"something_off": Noul(instructions="Is there any reason a warehouse lead would NOT mark this order as out for delivery today (wrong day, no customer, empty lines, canceled, implausible total)?")}
    p = client.system_one(state=state, questions=q).answers["something_off"].noul
    return p, f"P(something off) = {p:.2f}"


def gate(order, day):
    ok, reason = gate_deterministic(order)
    if not ok:
        return False, reason
    if GATE_MODE == "typesafe":
        p, note = gate_typesafe(order, day)
        if p is not None and p > GATE_ASK_ABOVE:
            return False, f"gate says ask a human: {note}"
        return True, f"{reason}; {note}"
    return True, reason


# ----------------------------------- the run -----------------------------------
def run(args):
    token = os.environ.get(TOKEN_ENV, "")
    after, before, day = today_range_utc()
    done = ledger_load()
    url = list_url(after, before)

    print(f"# status-agent {'DRY RUN' if args.dry_run else 'LIVE'} for {day}")
    print(f"# allowlist: {json.dumps({k: sorted(v) for k, v in ALLOWED_TRANSITIONS.items()})}  max writes: {MAX_WRITES_PER_RUN}  gate: {GATE_MODE}")
    print(f"# GET {url}")
    print(f"#   Authorization: Bearer {masked(token)}")

    if args.fixture:
        with open(args.fixture, encoding="utf-8") as f:
            orders = json.load(f)["data"]
        print(f"# fixture: {len(orders)} orders from {args.fixture} (no request sent)")
        reread = lambda oid: next((o for o in orders if o["id"] == oid), None)
    else:
        if not token:
            sys.exit(f"Set {TOKEN_ENV}, or pass --fixture sample-orders.json to rehearse offline.")
        orders = []
        while url:
            code, body = http("GET", url, token)
            if code != 200:
                sys.exit(f"GET failed {code}: {error_text(body)}")
            orders.extend(body.get("data", []))
            url = body.get("next_page")
        print(f"# {len(orders)} candidate order(s)")

        def reread(oid):
            code, body = http("GET", f"{BASE_URL}/orders/{oid}", token)
            return body.get("data") if code == 200 else None

    writes = 0
    counts = {"wrote": 0, "skipped": 0, "failed": 0, "stopped": 0}
    for o in orders:
        oid, num, status = o.get("id"), o.get("order_number"), o.get("status")
        base = {"order_id": oid, "order_number": num, "day": day, "from": status, "to": TARGET_STATUS}

        if writes >= MAX_WRITES_PER_RUN:
            ledger_write({**base, "action": "stopped", "reason": f"MAX_WRITES_PER_RUN={MAX_WRITES_PER_RUN} reached"}, args.dry_run)
            counts["stopped"] += 1
            continue
        if (oid, day, TARGET_STATUS) in done:
            ledger_write({**base, "action": "skipped", "reason": "already in ledger for today (idempotent)"}, args.dry_run)
            counts["skipped"] += 1
            continue

        # find-before-write: the list can lag a write by ~1 s, and a rep may have moved the order since
        fresh = reread(oid) or o
        status = fresh.get("status")
        base["from"] = status
        if TARGET_STATUS not in ALLOWED_TRANSITIONS.get(status, set()):
            ledger_write({**base, "action": "skipped", "reason": f"{status} -> {TARGET_STATUS} not in allowlist"}, args.dry_run)
            counts["skipped"] += 1
            continue

        ok, reason = gate(fresh, day)
        if not ok:
            ledger_write({**base, "action": "skipped", "reason": reason}, args.dry_run)
            counts["skipped"] += 1
            continue

        payload = {"id": oid, "status": TARGET_STATUS}
        if args.dry_run:
            ledger_write({**base, "action": "wrote", "reason": reason, "request": f"POST {BASE_URL}/orders {json.dumps(payload)}"}, dry_run=True)
            writes += 1
            counts["wrote"] += 1
            continue

        code, body = http("POST", f"{BASE_URL}/orders", token, payload)
        if code in (200, 201):
            ledger_write({**base, "action": "wrote", "reason": reason, "after": (body.get("data") or {}).get("status")}, dry_run=False)
            writes += 1
            counts["wrote"] += 1
            time.sleep(WRITE_PAUSE_SECONDS)
        elif code == 400:
            ledger_write({**base, "action": "failed", "reason": f"400 {error_text(body)}"}, dry_run=False)
            counts["failed"] += 1
        else:
            ledger_write({**base, "action": "failed", "reason": f"{code} {error_text(body)}"}, dry_run=False)
            counts["failed"] += 1
            print(f"# unexpected {code}; stopping this run so a human can look", file=sys.stderr)
            break

    print(f"# summary {day}: wrote={counts['wrote']} skipped={counts['skipped']} failed={counts['failed']} stopped={counts['stopped']}"
          + ("  (dry run: nothing was sent)" if args.dry_run else f"  ledger: {LEDGER_PATH}"))


def main():
    ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
    ap.add_argument("--dry-run", action="store_true", help="decide everything, write nothing")
    ap.add_argument("--fixture", help="JSON file shaped like GET /orders ({\"data\": [...]}) to rehearse without a key")
    args = ap.parse_args()
    if args.fixture and not args.dry_run:
        sys.exit("--fixture only makes sense with --dry-run")
    run(args)


if __name__ == "__main__":
    main()
