1#!/usr/bin/env python3
2"""Interrupt owned offline-cache writers and verify every retained intent and image."""
3import argparse
4import cache_images
5import hashlib
6import json
7import os
8from pathlib import Path
9import select
10import signal
11import subprocess
12import time
13
14
15def 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
120if __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())