| 1 | #!/usr/bin/env python3 |
| 2 | """Interrupt owned offline-cache writers and verify every retained intent and image.""" |
| 3 | import argparse |
| 4 | import cache_images |
| 5 | import hashlib |
| 6 | import json |
| 7 | import os |
| 8 | from pathlib import Path |
| 9 | import select |
| 10 | import signal |
| 11 | import subprocess |
| 12 | import time |
| 13 | |
| 14 | |
| 15 | def run(binary, output): |
| 16 | output.mkdir(parents=True, exist_ok=False) |
| 17 | cache = output / "cache.sqlite" |
| 18 | source = os.environ.get("ONESTORE_CACHE_PROBE_SOURCE") |
| 19 | source_hash = hashlib.sha256(Path(source).read_bytes()).hexdigest() if source else None |
| 20 | payload_bytes = int(os.environ.get("ONESTORE_CACHE_PROBE_BYTES", 2 * 1024 * 1024)) |
| 21 | assert 0 < payload_bytes <= 2 * 1024 * 1024 |
| 22 | (output / "run.json").write_text(json.dumps({ |
| 23 | "binary": str(binary), |
| 24 | "binary_sha256": hashlib.sha256(binary.read_bytes()).hexdigest(), |
| 25 | "controller_sha256": hashlib.sha256(Path(__file__).read_bytes()).hexdigest(), |
| 26 | "payload_bytes": payload_bytes, |
| 27 | "operation": os.environ.get("ONESTORE_CACHE_PROBE_OPERATION", "text"), |
| 28 | "source": source, |
| 29 | "source_sha256": source_hash, |
| 30 | }, indent=2)) |
| 31 | subprocess.run([binary, "init", cache], check=True, timeout=30) |
| 32 | retained = [] |
| 33 | retained_ids = [] |
| 34 | results = [] |
| 35 | delays = [None, "ack", "unack", "wal", 0, .001, .01, .02, .04, .08, .16, .32, .64, 1.28, 2.56, None, "ack"] |
| 36 | for operation, delay in enumerate(delays, 1): |
| 37 | with (output / f"owner-{operation}.stderr").open("w") as stderr: |
| 38 | owner = subprocess.Popen([binary, "edit", cache], stdin=subprocess.PIPE, |
| 39 | stdout=subprocess.PIPE, stderr=stderr, text=True, bufsize=1) |
| 40 | try: |
| 41 | def expect(prefix): |
| 42 | if not select.select([owner.stdout], [], [], 60)[0]: |
| 43 | raise TimeoutError(prefix) |
| 44 | actual = owner.stdout.readline().strip() |
| 45 | assert actual == prefix or actual.startswith(prefix + " "), (prefix, actual, owner.poll()) |
| 46 | return actual |
| 47 | |
| 48 | expect("ready") |
| 49 | contender = subprocess.run([binary, "read", cache], capture_output=True, text=True, timeout=15) |
| 50 | (output / f"contender-{operation}.stderr").write_text(contender.stderr) |
| 51 | assert contender.returncode != 0 and "DatabaseBusy" in contender.stderr, contender |
| 52 | owner.stdin.write(f"{operation} {'unack' if delay == 'unack' else 'ack'}\n") |
| 53 | owner.stdin.flush() |
| 54 | expect(f"editing {operation}") |
| 55 | acknowledgement = None |
| 56 | log_observation = None |
| 57 | if delay == "ack": |
| 58 | acknowledgement = expect(f"ack {operation}") |
| 59 | elif delay == "unack": |
| 60 | expect(f"durable {operation}") |
| 61 | elif delay == "wal": |
| 62 | # Cut the edit's transaction once its frames are in the log and its |
| 63 | # commit frame is not. |
| 64 | deadline = time.monotonic() + 10 |
| 65 | committed, _ = cache_images.wal_commits(cache) |
| 66 | while time.monotonic() < deadline: |
| 67 | commits, pending = cache_images.wal_commits(cache) |
| 68 | if commits == committed and pending: |
| 69 | log_observation = {"wal_commits": commits, "uncommitted_frames": pending} |
| 70 | break |
| 71 | if commits != committed: |
| 72 | break |
| 73 | assert log_observation, "No uncommitted log frames observed" |
| 74 | elif delay is not None: |
| 75 | time.sleep(delay) |
| 76 | owner.kill() |
| 77 | assert owner.wait(timeout=10) == -signal.SIGKILL, "Writer did not terminate at the requested process cut" |
| 78 | extra = owner.stdout.read() |
| 79 | if f"ack {operation} " in extra: |
| 80 | acknowledgement = extra.strip() |
| 81 | read = subprocess.run([binary, "read", cache], check=True, capture_output=True, |
| 82 | text=True, timeout=60) |
| 83 | actual = json.loads(read.stdout) |
| 84 | operations = actual["operations"] |
| 85 | ids = actual["ids"] |
| 86 | assert operations in (retained, retained + [operation]), (retained, operation, actual) |
| 87 | assert ids[:len(retained_ids)] == retained_ids, (retained_ids, actual) |
| 88 | if acknowledgement: |
| 89 | assert operations[-1] == operation |
| 90 | assert ids[-1] == int(acknowledgement.split()[2]) |
| 91 | if delay == "unack": |
| 92 | assert operations[-1] == operation and acknowledgement is None |
| 93 | assert actual["complete_payloads"] |
| 94 | results.append({"operation": operation, "delay": delay, |
| 95 | "process_exit": owner.returncode, |
| 96 | "acknowledged": acknowledgement is not None, |
| 97 | "retained": operation in operations, |
| 98 | "retained_operations": operations, "ids": ids, |
| 99 | "section_bytes": actual["section_bytes"], |
| 100 | "complete_payloads": True, "exclusive_owner": True, |
| 101 | "log_observation": log_observation}) |
| 102 | retained, retained_ids = operations, ids |
| 103 | (output / "results.json").write_text(json.dumps(results, indent=2)) |
| 104 | print(json.dumps(results[-1]), flush=True) |
| 105 | finally: |
| 106 | if owner.poll() is None: |
| 107 | owner.kill() |
| 108 | owner.wait(timeout=10) |
| 109 | owner.stdin.close() |
| 110 | owner.stdout.close() |
| 111 | assert any(row["acknowledged"] for row in results) |
| 112 | assert any(row["retained"] and not row["acknowledged"] for row in results) |
| 113 | assert any(not row["retained"] for row in results) |
| 114 | subprocess.run([binary, "read", cache, output / "recovered.one"], check=True, timeout=60) |
| 115 | if source: |
| 116 | assert hashlib.sha256(Path(source).read_bytes()).hexdigest() == source_hash |
| 117 | print(f"Passed {len(results)} offline-cache process interruptions", flush=True) |
| 118 | |
| 119 | |
| 120 | if __name__ == "__main__": |
| 121 | parser = argparse.ArgumentParser(description=__doc__) |
| 122 | parser.add_argument("output", type=Path) |
| 123 | parser.add_argument("--binary", type=Path, default=Path(__file__).resolve().parents[1] / "target/release/examples/cache_probe") |
| 124 | args = parser.parse_args() |
| 125 | run(args.binary.resolve(), args.output.resolve()) |