#!/usr/bin/env python3
"""Deterministic local failure trace. No network calls or real queue acknowledgements."""

from copy import deepcopy
from dataclasses import asdict
import json
from pathlib import Path
import subprocess
import sys
from tempfile import TemporaryDirectory

from consumer import IdentityConflict, ListingKey, ReadResult, Store, WrongKey, refresh_once


HERE = Path(__file__).resolve().parent


def main():
    fixture = json.loads((HERE / "fixtures/scenario.json").read_text())
    key = ListingKey.from_event(fixture["initial_event"])
    trace = []
    with TemporaryDirectory(prefix="listing-notification-lab-") as temporary:
        database = Path(temporary) / "inbox.sqlite3"
        # Hard process exit after COMMIT, with no close and no external queue ack.
        child = subprocess.run([
            sys.executable, "-B", "-c",
            "import os,sys; from consumer import Store; "
            "Store(sys.argv[1]).ingest(sys.stdin.read()); os._exit(23)",
            str(database),
        ], input=json.dumps(fixture["initial_event"]), text=True,
           cwd=HERE, capture_output=True, timeout=10)
        if child.returncode != 23:
            raise RuntimeError(f"unexpected child exit: {child.returncode}: {child.stderr}")
        with Store(database) as store:
            trace.append({"case": "exit_after_commit_before_queue_ack", "child_exit": 23,
                          "persisted_inbox": store.inbox_count(), "state": store.state(key)})
            duplicate = store.ingest(json.dumps(fixture["initial_event"]))
            trace.append({"case": "redelivery_after_restart", **asdict(duplicate),
                          "state": store.state(key)})
            initial = refresh_once(store, key, lambda k: ReadResult(k, fixture["initial_read"]))
            trace.append({"case": "initial_synthetic_read", "result": initial,
                          "state": store.state(key)})
            store.ingest(json.dumps(fixture["late_event"]))
            trace.append({"case": "older_event_is_a_new_trigger", "state": store.state(key)})
            ticket = store.begin_refresh(key)
            with Store(database) as ingester:
                ingester.ingest(json.dumps(fixture["during_read_event"]))
            applied = store.complete(ticket, ReadResult(key, fixture["discarded_read"]))
            trace.append({"case": "notification_during_read", "ticket_generation": ticket.generation,
                          "completion_applied": applied, "state": store.state(key)})

            def failing_read(_):
                raise OSError("synthetic read failure")

            try:
                refresh_once(store, key, failing_read)
            except OSError as error:
                trace.append({"case": "read_failure", "error": str(error), "state": store.state(key)})
            market_event = deepcopy(fixture["initial_event"])
            market_event.update(notification_id="notice-004", marketplace_id="market-B")
            seller_event = deepcopy(fixture["initial_event"])
            seller_event.update(notification_id="notice-005", seller_id="seller-B")
            for event in (market_event, seller_event):
                store.ingest(json.dumps(event))
            other_key = ListingKey.from_event(market_event)
            try:
                store.complete(store.begin_refresh(key), ReadResult(other_key, fixture["retry_read"]))
            except WrongKey:
                trace.append({"case": "wrong_key_result", "rejected": True,
                              "dirty_keys": [asdict(k) for k in store.dirty_keys()]})
            for pending in store.dirty_keys():
                refresh_once(store, pending, lambda k: ReadResult(k, fixture["retry_read"]))
            conflict = deepcopy(fixture["initial_event"])
            conflict["payload"] = {"hint": "different content under the same identity"}
            try:
                store.ingest(json.dumps(conflict))
            except IdentityConflict:
                trace.append({"case": "conflicting_identity", "rejected": True, "may_ack": False})
            final = {
                "inbox_rows": store.inbox_count(), "dirty_keys": len(store.dirty_keys()),
                "states": [{"key": asdict(k), **store.state(k)} for k in
                           (key, other_key, ListingKey.from_event(seller_event))],
            }
    print(json.dumps({
        "mode": "synthetic_local_only", "external_queue_ack_calls": 0,
        "trace": trace, "final": final,
    }, indent=2, sort_keys=True))


if __name__ == "__main__":
    main()
