jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204"""Real-process retry worker, host gate, zero-default, and backoff receipt."""
import http.serverimport osimport pathlibimport signalimport subprocessimport tempfileimport threadingimport timeimport urllib.errorimport urllib.request
ROOT = pathlib.Path(__file__).resolve().parents[1]RELAY_PORT = int(os.environ.get("STREAM_RETRY_CONFIG_RELAY_PORT", "6025"))STREAM_PORT = int(os.environ.get("STREAM_RETRY_CONFIG_PORT", "6026"))
class State: lock = threading.Lock() active = 0 peak = 0 requests = 0 first_at = None
@classmethod def reset(cls): with cls.lock: assert cls.active == 0 cls.peak = 0 cls.requests = 0 cls.first_at = None
class Handler(http.server.BaseHTTPRequestHandler): def do_GET(self): if self.path.startswith("/xrpc/com.atproto.sync.getRepo?"): with State.lock: State.active += 1 State.requests += 1 if State.first_at is None: State.first_at = time.monotonic() State.peak = max(State.peak, State.active) try: time.sleep(0.25) body = b'{"error":"ServiceUnavailable","message":"retry receipt"}' self.send_response(503) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) finally: with State.lock: State.active -= 1 return self.send_response(503) self.send_header("Content-Length", "0") self.end_headers()
def log_message(self, *_): pass
def run(*args): subprocess.run(args, cwd=ROOT, check=True)
def run_case(name, mode, count, workers, host_workers, expected_peak, min_peak=None): State.reset() with tempfile.TemporaryDirectory(prefix=f"stream-retry-{name}-") as raw: data = pathlib.Path(raw) / "data" log_path = pathlib.Path(raw) / "stream.log" data.mkdir() run("zig", "build", "write-sample", "--", "--retry-status", str(data), mode, str(count)) args = [ str(ROOT / "zig-out/bin/stream"), f"--port={STREAM_PORT}", f"--data-dir={data}", f"--upstream=ws://127.0.0.1:{RELAY_PORT}", f"--relay-http=http://127.0.0.1:{RELAY_PORT}", f"--plc=http://127.0.0.1:{RELAY_PORT}", "--compaction-interval=0", "--no-verify", "--failed-repo-retry-interval=1s", f"--failed-repo-retry-workers={workers}", f"--failed-repo-retry-host-workers={host_workers}", "--failed-repo-retry-max-delay=1500ms", ] with log_path.open("wb") as log: started_at = time.monotonic() process = subprocess.Popen(args, cwd=ROOT, stdout=log, stderr=subprocess.STDOUT) try: deadline = time.monotonic() + 20 while time.monotonic() < deadline: if process.poll() is not None: raise AssertionError(log_path.read_text(errors="replace")) text = log_path.read_text(errors="replace") if f"failed-repo retry pass: {count} candidates" in text: break time.sleep(0.05) else: raise AssertionError(log_path.read_text(errors="replace")) with State.lock: assert State.requests == count, (name, State.requests, text) # The cap is a ceiling, not a promise of simultaneity: reaching # it requires every worker to be in flight at once, and the # retry workers share an Io thread pool with the firehose # reconnect and accept loops. A small cap is saturated # regardless, so those cases still assert equality; the # 16-worker default was observed at 13, 14, 14, 16 and 16 across # runs on one idle 32-core box, so the overlap is scheduling # noise rather than a property of the host. Assert the ceiling # holds and that the default stays far above the explicit caps # -- pinning the exact overlap only asserts the scheduler. assert State.peak <= expected_peak, (name, State.peak, text) assert State.peak >= (expected_peak if min_peak is None else min_peak), ( name, State.peak, text, ) assert State.first_at is not None assert State.first_at - started_at >= 0.9, (name, State.first_at - started_at) finally: if process.poll() is None: process.send_signal(signal.SIGTERM) process.wait(timeout=10) run( "zig", "build", "write-sample", "--", "--retry-assert", str(data), str(count), "1000000", "1500000", ) print(f"retry config: PASS ({name}, peak={expected_peak})")
def run_disabled(): State.reset() with tempfile.TemporaryDirectory(prefix="stream-retry-disabled-") as raw: data = pathlib.Path(raw) / "data" log_path = pathlib.Path(raw) / "stream.log" data.mkdir() run("zig", "build", "write-sample", "--", "--retry-status", str(data), "unique", "6") args = [ str(ROOT / "zig-out/bin/stream"), f"--port={STREAM_PORT}", f"--data-dir={data}", f"--upstream=ws://127.0.0.1:{RELAY_PORT}", f"--relay-http=http://127.0.0.1:{RELAY_PORT}", f"--plc=http://127.0.0.1:{RELAY_PORT}", "--compaction-interval=0", "--no-verify", "--failed-repo-retry-interval=0", ] with log_path.open("wb") as log: process = subprocess.Popen(args, cwd=ROOT, stdout=log, stderr=subprocess.STDOUT) try: deadline = time.monotonic() + 10 while time.monotonic() < deadline: if process.poll() is not None: raise AssertionError(log_path.read_text(errors="replace")) try: urllib.request.urlopen(f"http://127.0.0.1:{STREAM_PORT}/healthz", timeout=0.2) break except OSError: time.sleep(0.05) else: raise AssertionError("disabled retry process never became healthy") time.sleep(1.2) with State.lock: assert State.requests == 0, State.requests finally: if process.poll() is None: process.send_signal(signal.SIGTERM) process.wait(timeout=10) print("retry config: PASS (disabled, no getRepo requests)")
server = http.server.ThreadingHTTPServer(("127.0.0.1", RELAY_PORT), Handler)thread = threading.Thread(target=server.serve_forever, daemon=True)thread.start()try: for _ in range(100): try: urllib.request.urlopen(f"http://127.0.0.1:{RELAY_PORT}/ready", timeout=0.1) except urllib.error.HTTPError: break except OSError: time.sleep(0.01) run_case("explicit-global", "unique", 6, 3, 1, 3) run_case("explicit-host", "same", 6, 4, 2, 2) # min_peak 8 is twice the largest explicit cap below it (default-host's 4), # so a regression to any smaller worker default still fails loudly. run_case("default-global", "unique", 20, 0, 1, 16, min_peak=8) run_case("default-host", "same", 8, 8, 0, 4) run_disabled()finally: server.shutdown() server.server_close() thread.join(timeout=2)