#!/usr/bin/env python3
"""Load balancing + redelivery reorders a broker's "in-order" stream, DDIA 2e Ch. 12.

A producer sends five messages in order: m1, m2, m3, m4, m5. A message broker
that hands each message to a consumer and expects an ACK on completion is often
described as preserving delivery order. This script shows that two features the
broker offers -- LOAD BALANCING (fanning messages across several consumers for
throughput) and REDELIVERY (re-sending a message whose lease expires because a
consumer stalled) -- together break the overall order.

Scenario A (shared queue, two consumers, load balancing + redelivery):
    The broker fans the five messages across C1 and C2. C2 stalls mid-m3 (a GC
    pause). Its lease on m3 expires, so the broker REDELIVERS m3 to C1 -- but by
    then C1 has already processed m4. So m3 (sent 3rd) is ACKed AFTER m4 (sent
    4th). This is DDIA Figure 12-2.

Scenario B (single partitioned log, one consumer, same stall):
    One consumer owns the log and processes it in offset order. The SAME stall
    happens on m3, but there is no second consumer to redeliver to, so nothing
    overtakes m3: the consumer simply blocks on it (head-of-line blocking) and
    resumes. Completion order equals send order -- slower, but ordered.

Everything is modeled explicitly on a virtual clock, so it is deterministic and
reproducible. Pure standard library.
"""
import heapq
from collections import deque

SEND_ORDER = ["m1", "m2", "m3", "m4", "m5"]

TICK = 10       # ticks to process one message on a healthy consumer
LEASE = 8       # ticks the broker waits for an ACK before redelivering
PAUSE = 30      # ticks a consumer hangs (e.g. a GC pause) while holding a message

# The fault: C2 hangs while holding m3.
STALL_MSG = "m3"
STALL_CONSUMER = "C2"


def run_shared_queue():
    """Two consumers, load-balanced over a shared queue, with lease-timeout
    redelivery. Returns (event_log, completion_order).

    The broker's load-balanced fan-out (prefetch) assigns:
        C1 <- m1, m4, m5      C2 <- m2, m3
    C2 hangs on m3; its lease expires and the broker redelivers m3 to C1, where
    it is inserted at the front of C1's remaining work (a redelivery has waited
    longest, so it jumps the queue). We compute when everything actually
    completes -- the ordering is an OUTPUT of the simulation, not hard-coded.
    """
    log = []
    completions = []          # (time, msg) in the order ACKs actually land
    events = []               # heap of (time, seq, kind, data)
    seq = 0

    def sched(time, kind, **data):
        nonlocal seq
        heapq.heappush(events, (time, seq, kind, data))
        seq += 1

    queues = {"C1": deque(["m1", "m4", "m5"]),
              "C2": deque(["m2", "m3"])}
    busy = {"C1": False, "C2": False}

    def start(consumer, msg, t):
        busy[consumer] = True
        if msg == STALL_MSG and consumer == STALL_CONSUMER:
            log.append((t, f"{consumer} starts {msg}, then HANGS (GC pause) -- no ACK coming"))
            # The broker's lease fires first; the consumer only wakes at t+PAUSE,
            # by which point its lease is long gone.
            sched(t + LEASE, "lease_expire", msg=msg, consumer=consumer)
            sched(t + PAUSE, "complete", msg=msg, consumer=consumer, stalled=True)
        else:
            log.append((t, f"{consumer} starts {msg}"))
            sched(t + TICK, "complete", msg=msg, consumer=consumer, stalled=False)

    # Kick off: each consumer begins its first prefetched message at t=0.
    for c in ("C1", "C2"):
        start(c, queues[c].popleft(), 0)

    while events:
        t, _, kind, d = heapq.heappop(events)

        if kind == "complete":
            c, msg = d["consumer"], d["msg"]
            if d["stalled"]:
                # The hung consumer finally wakes, but its lease already expired
                # and the broker gave the message to someone else.
                log.append((t, f"{c} wakes and finishes {msg}, but its lease is gone -> ACK rejected (duplicate)"))
                busy[c] = False
                continue
            log.append((t, f"{c} ACKs {msg}  [DONE]"))
            completions.append((t, msg))
            busy[c] = False
            if queues[c]:
                start(c, queues[c].popleft(), t)

        elif kind == "lease_expire":
            msg, holder = d["msg"], d["consumer"]
            log.append((t, f"broker: lease on {msg} (held by {holder}) EXPIRED -> redeliver to C1"))
            # Redelivery jumps the front of the healthy consumer's queue.
            queues["C1"].appendleft(msg)
            if not busy["C1"]:
                start("C1", queues["C1"].popleft(), t)

    return log, completions


def run_partitioned_log():
    """One consumer owning a single partitioned log, processed in offset order.
    The SAME stall happens on m3, but there is no competitor to redeliver to, so
    the consumer just blocks on m3 and resumes. Returns (event_log, completion_order).
    """
    log = []
    completions = []
    t = 0
    for msg in SEND_ORDER:
        if msg == STALL_MSG:
            log.append((t, f"C1 starts {msg}, then HANGS (GC pause, {PAUSE} ticks)"))
            log.append((t, "  no second consumer exists -> broker cannot redeliver; nothing overtakes m3"))
            t += PAUSE
        else:
            log.append((t, f"C1 starts {msg}"))
            t += TICK
        log.append((t, f"C1 ACKs {msg}  [DONE]"))
        completions.append((t, msg))
    return log, completions


def print_log(log):
    for t, line in log:
        print(f"    [t={t:>2}] {line}")


def order_str(completions):
    return ", ".join(msg for _, msg in completions)


def describe_reorder(completions):
    """Find and report the send-order inversion, if any."""
    order = [msg for _, msg in completions]
    send_index = {m: i for i, m in enumerate(SEND_ORDER)}
    for i in range(len(order)):
        for j in range(i + 1, len(order)):
            if send_index[order[i]] > send_index[order[j]]:
                return (f"{order[j]} was sent BEFORE {order[i]}, but completed AFTER it "
                        f"-- delivery order was NOT preserved.")
    return "completion order == send order -- delivery order preserved."


def main():
    print("A broker sends five messages in order, expecting an ACK per message.")
    print(f"    send order: {', '.join(SEND_ORDER)}\n")

    print("=" * 68)
    print("Scenario A: shared queue, load balancing across C1 + C2 + redelivery")
    print("=" * 68)
    print("Broker fans out (prefetch):  C1 <- m1, m4, m5   |   C2 <- m2, m3")
    print(f"Lease = {LEASE} ticks; a healthy message takes {TICK}; a hung consumer stalls {PAUSE}.\n")
    log_a, comp_a = run_shared_queue()
    print_log(log_a)
    print(f"\n    send order:       {', '.join(SEND_ORDER)}")
    print(f"    completion order: {order_str(comp_a)}")
    print(f"    -> {describe_reorder(comp_a)}\n")

    print("=" * 68)
    print("Scenario B: single partitioned log, one consumer, SAME stall on m3")
    print("=" * 68)
    print("No load balancing: one consumer processes the log in offset order.\n")
    log_b, comp_b = run_partitioned_log()
    print_log(log_b)
    print(f"\n    send order:       {', '.join(SEND_ORDER)}")
    print(f"    completion order: {order_str(comp_b)}")
    print(f"    -> {describe_reorder(comp_b)}")


if __name__ == "__main__":
    main()
