|
| 1 | +from __future__ import annotations |
| 2 | + |
| 3 | +import json |
| 4 | +import time |
| 5 | +from dataclasses import dataclass |
| 6 | + |
| 7 | +import httpx |
| 8 | + |
| 9 | +from pypproxy.store.models import Entry |
| 10 | + |
| 11 | + |
| 12 | +@dataclass |
| 13 | +class ABResult: |
| 14 | + endpoint_a: str |
| 15 | + endpoint_b: str |
| 16 | + method: str |
| 17 | + status_a: int |
| 18 | + status_b: int |
| 19 | + body_a: bytes |
| 20 | + body_b: bytes |
| 21 | + duration_a_ms: int |
| 22 | + duration_b_ms: int |
| 23 | + error_a: str = "" |
| 24 | + error_b: str = "" |
| 25 | + |
| 26 | + @property |
| 27 | + def status_diff(self) -> bool: |
| 28 | + return self.status_a != self.status_b |
| 29 | + |
| 30 | + @property |
| 31 | + def body_diff(self) -> bool: |
| 32 | + return self.body_a != self.body_b |
| 33 | + |
| 34 | + def to_dict(self) -> dict: |
| 35 | + import base64 |
| 36 | + |
| 37 | + return { |
| 38 | + "endpoint_a": self.endpoint_a, |
| 39 | + "endpoint_b": self.endpoint_b, |
| 40 | + "method": self.method, |
| 41 | + "status_a": self.status_a, |
| 42 | + "status_b": self.status_b, |
| 43 | + "body_a": base64.b64encode(self.body_a).decode() if self.body_a else "", |
| 44 | + "body_b": base64.b64encode(self.body_b).decode() if self.body_b else "", |
| 45 | + "duration_a_ms": self.duration_a_ms, |
| 46 | + "duration_b_ms": self.duration_b_ms, |
| 47 | + "error_a": self.error_a, |
| 48 | + "error_b": self.error_b, |
| 49 | + "status_diff": self.status_diff, |
| 50 | + "body_diff": self.body_diff, |
| 51 | + } |
| 52 | + |
| 53 | + def diff_summary(self) -> str: |
| 54 | + lines: list[str] = [] |
| 55 | + if self.status_diff: |
| 56 | + lines.append(f"Status: A={self.status_a} B={self.status_b}") |
| 57 | + else: |
| 58 | + lines.append(f"Status: both {self.status_a}") |
| 59 | + if self.body_diff: |
| 60 | + lines.append(f"Body differs ({len(self.body_a):,} B vs {len(self.body_b):,} B)") |
| 61 | + # Try JSON diff summary |
| 62 | + try: |
| 63 | + da = json.loads(self.body_a) |
| 64 | + db = json.loads(self.body_b) |
| 65 | + if isinstance(da, dict) and isinstance(db, dict): |
| 66 | + added = set(db) - set(da) |
| 67 | + removed = set(da) - set(db) |
| 68 | + changed = {k for k in da if k in db and da[k] != db[k]} |
| 69 | + if added: |
| 70 | + lines.append(f" + fields added: {', '.join(sorted(added)[:5])}") |
| 71 | + if removed: |
| 72 | + lines.append(f" - fields removed: {', '.join(sorted(removed)[:5])}") |
| 73 | + if changed: |
| 74 | + lines.append(f" ~ fields changed: {', '.join(sorted(changed)[:5])}") |
| 75 | + except Exception: |
| 76 | + pass |
| 77 | + else: |
| 78 | + lines.append("Body: identical") |
| 79 | + lines.append(f"Latency: A={self.duration_a_ms}ms B={self.duration_b_ms}ms") |
| 80 | + return "\n".join(lines) |
| 81 | + |
| 82 | + |
| 83 | +async def run_ab_test( |
| 84 | + entry: Entry, |
| 85 | + override_host_b: str, |
| 86 | + override_scheme_b: str = "", |
| 87 | + timeout: int = 30, |
| 88 | +) -> ABResult: |
| 89 | + """Send the same request to two different hosts and compare responses.""" |
| 90 | + scheme = entry.scheme |
| 91 | + path = entry.path |
| 92 | + query = entry.query |
| 93 | + headers = { |
| 94 | + k: ", ".join(v) |
| 95 | + for k, v in entry.req_headers.items() |
| 96 | + if k.lower() not in ("host", "content-length", "connection") |
| 97 | + } |
| 98 | + body = entry.req_body |
| 99 | + |
| 100 | + url_a = f"{scheme}://{entry.host}{path}" + (f"?{query}" if query else "") |
| 101 | + scheme_b = override_scheme_b or scheme |
| 102 | + url_b = f"{scheme_b}://{override_host_b}{path}" + (f"?{query}" if query else "") |
| 103 | + |
| 104 | + status_a = status_b = 0 |
| 105 | + body_a = body_b = b"" |
| 106 | + dur_a = dur_b = 0 |
| 107 | + err_a = err_b = "" |
| 108 | + |
| 109 | + async def _fetch(url: str) -> tuple[int, bytes, int, str]: |
| 110 | + start = time.monotonic() |
| 111 | + try: |
| 112 | + h = dict(headers) |
| 113 | + h["host"] = url.split("/")[2].split(":")[0] |
| 114 | + async with httpx.AsyncClient(verify=False, timeout=timeout, http2=True) as client: |
| 115 | + resp = await client.request(method=entry.method, url=url, headers=h, content=body) |
| 116 | + return resp.status_code, resp.content, int((time.monotonic() - start) * 1000), "" |
| 117 | + except Exception as e: |
| 118 | + return 0, b"", int((time.monotonic() - start) * 1000), str(e) |
| 119 | + |
| 120 | + (status_a, body_a, dur_a, err_a), (status_b, body_b, dur_b, err_b) = await __import__( |
| 121 | + "asyncio" |
| 122 | + ).gather(_fetch(url_a), _fetch(url_b)) |
| 123 | + |
| 124 | + return ABResult( |
| 125 | + endpoint_a=url_a, |
| 126 | + endpoint_b=url_b, |
| 127 | + method=entry.method, |
| 128 | + status_a=status_a, |
| 129 | + status_b=status_b, |
| 130 | + body_a=body_a, |
| 131 | + body_b=body_b, |
| 132 | + duration_a_ms=dur_a, |
| 133 | + duration_b_ms=dur_b, |
| 134 | + error_a=err_a, |
| 135 | + error_b=err_b, |
| 136 | + ) |
0 commit comments