diff --git a/.github/workflows/pagestore-test.yml b/.github/workflows/pagestore-test.yml index 49913754ab35a..258f5ffab3d6c 100644 --- a/.github/workflows/pagestore-test.yml +++ b/.github/workflows/pagestore-test.yml @@ -76,13 +76,14 @@ jobs: cc -O2 -Wall -Wextra -Werror -o pagestore_daemon pagestore_daemon.c pagestore_core.c pagestore_fault.c storage_posix.c pagestore_layer.c pagestore_layer_store.c pagestore_manifest.c pagestore_memtable.c pagestore_pgcache.c pagestore_prune.c pagestore_retention.c pagestore_wal_store.c pagestore_wal_segment.c pagestore_walidx_prune.c pagestore_walidx_snapshot.c pagestore_forkmeta_prune.c pagestore_forkmeta_snapshot.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_inspect pagestore_inspect.c -lrt cc -O2 -Wall -Wextra -Werror -o pagestore_test pagestore_test.c -lrt - cc -O2 -Wall -Wextra -Werror -o pagestore_layer_test pagestore_layer_test.c pagestore_layer.c pagestore_layer_store.c -lrt - cc -O2 -Wall -Wextra -Werror -o pagestore_layer_store_test pagestore_layer_store_test.c pagestore_layer_store.c pagestore_layer.c -lrt + cc -O2 -Wall -Wextra -Werror -o pagestore_layer_test pagestore_layer_test.c pagestore_layer.c pagestore_layer_store.c pagestore_fault.c -lrt -pthread + cc -O2 -Wall -Wextra -Werror -o pagestore_layer_store_test pagestore_layer_store_test.c pagestore_layer_store.c pagestore_layer.c pagestore_fault.c -lrt -pthread + cc -O2 -Wall -Wextra -Werror -o pagestore_layer_crash_client pagestore_layer_crash_client.c -lrt -pthread cc -O2 -Wall -Wextra -Werror -o pagestore_tiering_test pagestore_tiering_test.c pagestore_core.c pagestore_fault.c storage_posix.c pagestore_layer.c pagestore_layer_store.c pagestore_manifest.c pagestore_memtable.c pagestore_pgcache.c pagestore_prune.c pagestore_retention.c pagestore_wal_store.c pagestore_wal_segment.c pagestore_walidx_prune.c pagestore_walidx_snapshot.c pagestore_forkmeta_prune.c pagestore_forkmeta_snapshot.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_forkmeta_cutover_test pagestore_forkmeta_cutover_test.c pagestore_core.c pagestore_fault.c storage_posix.c pagestore_layer.c pagestore_layer_store.c pagestore_manifest.c pagestore_memtable.c pagestore_pgcache.c pagestore_prune.c pagestore_retention.c pagestore_wal_store.c pagestore_wal_segment.c pagestore_walidx_prune.c pagestore_walidx_snapshot.c pagestore_forkmeta_prune.c pagestore_forkmeta_snapshot.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_forkmeta_crash_matrix_test pagestore_forkmeta_crash_matrix_test.c pagestore_core.c pagestore_fault.c storage_posix.c pagestore_layer.c pagestore_layer_store.c pagestore_manifest.c pagestore_memtable.c pagestore_pgcache.c pagestore_prune.c pagestore_retention.c pagestore_wal_store.c pagestore_wal_segment.c pagestore_walidx_prune.c pagestore_walidx_snapshot.c pagestore_forkmeta_prune.c pagestore_forkmeta_snapshot.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_pgcache_test pagestore_pgcache_test.c pagestore_pgcache.c -lrt -lpthread - cc -O2 -Wall -Wextra -Werror -o pagestore_manifest_test pagestore_manifest_test.c pagestore_manifest.c pagestore_layer.c pagestore_layer_store.c -lrt + cc -O2 -Wall -Wextra -Werror -o pagestore_manifest_test pagestore_manifest_test.c pagestore_manifest.c pagestore_layer.c pagestore_layer_store.c pagestore_fault.c -lrt -pthread cc -O2 -Wall -Wextra -Werror -o pagestore_retention_test pagestore_retention_test.c pagestore_retention.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_gc_test pagestore_gc_test.c pagestore_core.c pagestore_fault.c storage_posix.c pagestore_layer.c pagestore_layer_store.c pagestore_manifest.c pagestore_memtable.c pagestore_pgcache.c pagestore_prune.c pagestore_retention.c pagestore_wal_store.c pagestore_wal_segment.c pagestore_walidx_prune.c pagestore_walidx_snapshot.c pagestore_forkmeta_prune.c pagestore_forkmeta_snapshot.c -lrt -lpthread cc -O2 -Wall -Wextra -Werror -o pagestore_walidx_prune_test pagestore_walidx_prune_test.c pagestore_walidx_prune.c -lrt @@ -109,6 +110,22 @@ jobs: --daemon-binary ./pagestore_daemon \ --inspect-binary ./pagestore_inspect + - name: Run image-layer H1 crash scenarios + working-directory: contrib/pagestore + run: | + for scenario in \ + image_layer_after_create \ + image_layer_after_write \ + image_layer_after_seal \ + image_layer_after_manifest_add; do + python3 harness/pagestore_harness.py \ + --capabilities harness/capabilities.json \ + --daemon-fault-recovery "harness/scenarios/${scenario}.jsonl" \ + --daemon-binary ./pagestore_daemon \ + --layer-client-binary ./pagestore_layer_crash_client \ + --inspect-binary ./pagestore_inspect + done + - name: Run image-layer unit test working-directory: contrib/pagestore run: ./pagestore_layer_test diff --git a/contrib/pagestore/MVP_COMPLETION_PLAN.md b/contrib/pagestore/MVP_COMPLETION_PLAN.md index 7c53111378b0b..9867d81e75c6a 100644 --- a/contrib/pagestore/MVP_COMPLETION_PLAN.md +++ b/contrib/pagestore/MVP_COMPLETION_PLAN.md @@ -783,8 +783,9 @@ so the SPDK frontend remains runnable without that POSIX-only mailbox. The first composed H1 materializer slice is now implemented: the two restartpoint plans pause the checkpointer child after relation-page sync/before marker write and after marker sync, then stop and recover the whole -materializer. The remaining branch, layer, reclaim, and GC H1 families remain -outstanding. +materializer. The prepared-receipt/service-restore branch slice and the POSIX +image-layer create/write/seal/manifest-ADD publication slice are also covered. +Branch bootstrap/install, manifest replacement, reclaim, and GC H1 cases remain. Deliverables: @@ -810,8 +811,10 @@ Expected scope: one or two PRs. ### H1. Compose process-level crash scenarios -Status: **materializer replay/restartpoint slice implemented; branch, layer, -reclaim, and GC cases remain and depend on H0/R2-R5**. +Status: **materializer replay/restartpoint, branch prepared-receipt/service- +restore, and POSIX image-layer publication slices implemented; branch +bootstrap/install, manifest replacement, reclaim, and GC cases remain and +depend on H0/R2-R5**. Required scenario families: diff --git a/contrib/pagestore/MVP_STATUS.md b/contrib/pagestore/MVP_STATUS.md index a1eceb676b5da..e9507971b4298 100644 --- a/contrib/pagestore/MVP_STATUS.md +++ b/contrib/pagestore/MVP_STATUS.md @@ -42,7 +42,7 @@ new daemon's zeroing/recovery window. | Area | Status | Current proof | |---|---|---| | Page ingest and copy-on-write reads | Implemented | standalone and PostgreSQL integration suites | -| Image-layer path | Functional mechanisms implemented; phases 2–3 partial | manifest/compaction/segment-GC restart tests; sparse indexes and layer-block cache invalidation remain | +| Image-layer path | Functional mechanisms implemented; H1 POSIX publication crash slice covered; phases 2–3 partial | manifest/compaction/segment-GC restart tests plus create/write/seal/manifest-ADD crash recovery scenarios with sentinel/LSN and idempotent restart checks; sparse indexes and layer-block cache invalidation remain | | Filesystem object tier | Upload done; cache/GC operations partial | download, eviction, refresh, and remote-delete tests; cache policy and orphan reconciliation remain | | Materialized-page cache | Basic version cache implemented; phase partial | bounded cache/invalidation tests; cost-aware admission and integrated redo avoidance remain | | WAL shipping and ancestry-aware WAL reads | Immutable 1 MiB segments integrated for sealed prefixes; the flat-log copy of every complete sealed record is reclaimed, while the flat log remains migration/tail authority | chunk assembly, reopen, ancestry, and WAL segment/store tests | @@ -78,6 +78,18 @@ The existing CI proves both focused subsystem paths and the composed contract: - `branch_boot_test.sh` proves an independent branch compute can boot, preserve fork-point visibility, and write on its own timeline. +The H1 image-layer crash slice is limited to POSIX local layers and process +abort. Its four ordered stages cover canonical file creation, file writes +before seal, sealed layer data, and durable `layers.manifest` ADD publication; +the write stage does not claim power-loss durability. The composed +harness keeps the pre-recovery physical snapshot, then checks a sentinel +page/LSN, the expected manifest state, and one additional restart for +idempotence. A crash after ADD but before its flush watermark conservatively +retains the durable layer and republishes segment-backed coverage once; the +intervening clean shutdown compacts that conservative duplicate, and the second +restart proves convergence to one layer. The slice does not delete crash +orphans; cross-process ownership is required before adding that policy. + ## MVP gates ### 1. One composed golden scenario -- implemented @@ -308,12 +320,13 @@ mutations and its snapshot maintenance can be forced below the geometric trigger. All three controllers are POSIX-only; forkmeta does not claim the remaining R6 queue-bound soak/tuning work. -### 5. Composed crash and format-compatibility coverage -- remaining +### 5. Composed crash and format-compatibility coverage -- partial -Focused crash-safety tests exist, while the declarative harness remains partial. +The POSIX image-layer publication slice is now covered by the declarative +harness. Other crash boundaries remain outside this slice. Before declaring the MVP repeatable, add process-level fault scenarios around -materializer progress, branch prepare/install, manifest/layer publication, and -retention GC, plus a persisted-format fixture for restart/upgrade compatibility. +branch bootstrap/install, manifest replacement, and retention/reclaim/GC, plus +a persisted-format fixture for restart/upgrade compatibility. ## Recommended sequence diff --git a/contrib/pagestore/harness/capabilities.json b/contrib/pagestore/harness/capabilities.json index f3d82a4cd379e..4b16877b028f1 100644 --- a/contrib/pagestore/harness/capabilities.json +++ b/contrib/pagestore/harness/capabilities.json @@ -13,11 +13,12 @@ "io_unit": 262144 }, "daemon_fault_smoke": { - "operations": ["crash", "set_fault", "release_fault"], + "operations": ["crash", "set_fault", "release_fault", "layer_seed"], "constraints": { "crash": {}, "set_fault": {}, - "release_fault": {"target": ["store"]} + "release_fault": {"target": ["store"]}, + "layer_seed": {"target": ["store"]} }, "protocol_version": 45, "page_size": 8192, @@ -71,6 +72,7 @@ "install_reader", "restart", "crash", + "layer_seed", "advance", "assert", "materializer_fault", diff --git a/contrib/pagestore/harness/pagestore_harness.py b/contrib/pagestore/harness/pagestore_harness.py index af84004ae9d45..85ccc0a930cf6 100644 --- a/contrib/pagestore/harness/pagestore_harness.py +++ b/contrib/pagestore/harness/pagestore_harness.py @@ -272,6 +272,7 @@ def postgresql_conf_string(value: str | Path) -> str: "sync": {"op", "id", "target", "kind", "extra"}, "set_fault": {"op", "id", "target", "fault", "action", "hit", "timeout", "extra"}, "release_fault": {"op", "id", "target", "fault", "extra"}, + "layer_seed": {"op", "id", "target", "extra"}, "capture": {"op", "id", "target", "kind", "name", "horizon", "extra"}, "compare": {"op", "id", "left", "right", "extra"}, "expect_failure": {"op", "id", "target", "command", "sqlstate", "extra"}, @@ -298,6 +299,7 @@ def postgresql_conf_string(value: str | Path) -> str: "sync": {"target", "kind"}, "set_fault": {"target", "fault", "action"}, "release_fault": {"target", "fault"}, + "layer_seed": {"target"}, "capture": {"target", "kind", "name", "horizon"}, "compare": {"left", "right"}, "expect_failure": {"target", "command"}, @@ -645,7 +647,7 @@ def validate_fault_action( RUNTIME_OPERATIONS = { "daemon_smoke": {"crash"}, - "daemon_fault_smoke": {"crash", "set_fault", "release_fault"}, + "daemon_fault_smoke": {"crash", "set_fault", "release_fault", "layer_seed"}, "writer_smoke": { "sql", "checkpoint", "prepare_reader", "reader_base", "bootstrap", "install_reader", "assert", "capture", @@ -666,6 +668,7 @@ def validate_fault_action( "crash": {}, "set_fault": {}, "release_fault": {"target": ["store"]}, + "layer_seed": {"target": ["store"]}, }, "writer_smoke": { "checkpoint": {"target": ["writer"]}, @@ -964,6 +967,11 @@ def validate_runtime_plan(plan: Plan, capabilities: dict[str, Any], runtime: str f"runtime {runtime!r} uses materializer_fault for named process faults" ) elif runtime == "daemon_fault_smoke": + seed_actions = [action for action in plan.actions if action["op"] == "layer_seed"] + if len(seed_actions) > 1: + raise PlanError( + f"runtime {runtime!r} permits at most one layer_seed action" + ) named = [ action for action in plan.actions if action["op"] in ("crash", "set_fault") and "fault" in action @@ -975,7 +983,18 @@ def validate_runtime_plan(plan: Plan, capabilities: dict[str, Any], runtime: str raise PlanError( "runtime daemon_fault_smoke requires exactly one named fault action" ) + if seed_actions and plan.actions.index(seed_actions[0]) > plan.actions.index(named[0]): + raise PlanError( + "runtime daemon_fault_smoke requires layer_seed before the named fault" + ) fault = named[0] + if seed_actions and fault["fault"] not in { + "image_layer.after_create", "image_layer.after_write", + "image_layer.after_seal", "image_layer.after_manifest_add", + }: + raise PlanError( + "runtime daemon_fault_smoke layer_seed requires an H1 image-layer fault" + ) releases = [ action for action in plan.actions if action["op"] == "release_fault" and action["fault"] == fault["fault"] @@ -2203,6 +2222,143 @@ def _capture_fault_diagnostics( return root / "fault-diagnostics.json" +def _layer_fault_stage(fault_name: str) -> str | None: + stages = { + "image_layer.after_create": "create", + "image_layer.after_write": "write", + "image_layer.after_seal": "seal", + "image_layer.after_manifest_add": "manifest_add", + } + return stages.get(fault_name) + + +def _canonical_layer_files(store: Path) -> list[Path]: + files: list[Path] = [] + pattern = re.compile(r"^layer_(?:0|[1-9][0-9]*)_[0-9a-fA-F]{16}$") + for path in sorted(store.iterdir()): + if not pattern.fullmatch(path.name): + continue + value = path.lstat() + if stat.S_ISREG(value.st_mode): + files.append(path) + return files + + +def _check_layer_crash_snapshot(store: Path, stage: str) -> None: + """Check process-abort physical state before recovery mutates the store. + + The after_write point deliberately proves only the write-stage ordering; + it is not a power-loss or file-fsync durability claim. + """ + files = _canonical_layer_files(store) + if len(files) != 1: + raise OracleMismatch( + f"H1 {stage} crash left {len(files)} canonical layer files, expected one" + ) + manifest = store / "layers.manifest" + manifest_size = manifest.stat().st_size if manifest.exists() else 0 + size = files[0].stat().st_size + if stage == "create" and size != 0: + raise OracleMismatch("after_create did not leave an empty canonical layer") + if stage in {"write", "seal"} and size == 0: + raise OracleMismatch(f"after_{stage} did not leave written layer bytes") + if stage != "manifest_add" and manifest_size != 0: + raise OracleMismatch(f"after_{stage} unexpectedly published layers.manifest") + if stage == "manifest_add" and manifest_size == 0: + raise OracleMismatch("after_manifest_add did not leave a durable manifest ADD") + + +def _check_layer_manifest( + manifest: dict[str, Any], stage: str, recovery_restart: bool, +) -> None: + # Before manifest publication, segment replay rebuilds and republishes the + # interrupted flush. After ADD publication but before the flush watermark, + # recovery conservatively retains that layer and republishes segment-backed + # coverage once. The intervening clean shutdown compacts that conservative + # duplicate, so the following restart must converge to one layer. + expected_states = ( + ((1, 1), (2, 2)) + if stage == "manifest_add" and not recovery_restart + else ((1, 1),) + ) + expected_text = "(1, 1) or (2, 2)" if len(expected_states) > 1 else "(1, 1)" + observed_state = (manifest.get("layer_count"), manifest.get("local_layers")) + if observed_state not in expected_states: + raise OracleMismatch( + f"after_{stage} recovery reported layer_count/local_layers=" + f"{observed_state!r}, expected {expected_text}" + ) + if manifest.get("manifest_poisoned") is not False: + raise OracleMismatch( + f"after_{stage} recovery reported manifest_poisoned=" + f"{manifest.get('manifest_poisoned')!r}" + ) + + +def _check_layer_manifest_after_restart( + inspector: Path, + shm: str, + inspection_schema: dict[str, Any], + stage: str, + timeout: float, +) -> None: + """Wait for the second restart to compact a manifest_add duplicate.""" + poll_timeout = max(0.0, min(10.0, timeout)) + deadline = time.monotonic() + poll_timeout + while True: + manifest = inspect_store(inspector, shm, "manifest", inspection_schema) + observed_state = (manifest.get("layer_count"), manifest.get("local_layers")) + if manifest.get("manifest_poisoned") is not False: + raise OracleMismatch( + f"after_{stage} recovery reported manifest_poisoned=" + f"{manifest.get('manifest_poisoned')!r}" + ) + if observed_state == (1, 1): + _check_layer_manifest(manifest, stage, True) + return + + # Only the coherent duplicate produced by manifest_add recovery is + # allowed to remain transient. Poisoned or otherwise inconsistent + # observations must not be hidden by a later successful poll. + if ( + stage != "manifest_add" + or observed_state != (2, 2) + ): + _check_layer_manifest(manifest, stage, True) + + now = time.monotonic() + if now >= deadline: + raise HarnessTimeout( + f"after_{stage} recovery manifest did not converge to (1, 1) " + f"within {poll_timeout:.3f}s; last state={observed_state!r}" + ) + time.sleep(min(0.05, deadline - now)) + + +def _start_layer_client( + client: Path, shm: str, mode: str, log: Path, +) -> subprocess.Popen[str]: + with log.open("a", encoding="utf-8") as output: + return subprocess.Popen( + [str(client.resolve()), "--shm", shm, "--mode", mode], + stdout=output, stderr=subprocess.STDOUT, text=True, + env=private_environment(), start_new_session=True, + ) + + +def _verify_layer_client(client: Path, shm: str, log: Path, timeout: float) -> None: + with log.open("a", encoding="utf-8") as output: + result = subprocess.run( + [str(client.resolve()), "--shm", shm, "--mode", "verify"], + stdout=output, stderr=subprocess.STDOUT, text=True, + env=private_environment(), timeout=max(5.0, timeout), check=False, + ) + if result.returncode != 0: + raise OracleMismatch( + f"layer sentinel client failed after recovery with status {result.returncode}" + ) + + def run_daemon_fault_recovery( plan: Plan, capabilities: dict[str, Any], @@ -2214,10 +2370,12 @@ def run_daemon_fault_recovery( capabilities_path: Path | None = None, rerun_command: list[str] | None = None, timeout: float = 15.0, + layer_client: Path | None = None, ) -> Path: """Run one pre-armed named daemon fault and prove recovery is idempotent.""" daemon = daemon.resolve() inspector = inspector.resolve() + layer_client = layer_client.resolve() if layer_client is not None else None validate_plan(plan, capabilities, capabilities_path) validate_runtime_plan(plan, capabilities, "daemon_fault_smoke") root, temporary = run_root(requested_root) @@ -2260,6 +2418,10 @@ def run_daemon_fault_recovery( ): raise PlanError("daemon fault recovery requires exactly one named fault action") action = named_faults[0] + layer_seed_actions = [item for item in plan.actions if item["op"] == "layer_seed"] + if layer_seed_actions and layer_client is None: + raise PlanError("layer_seed requires --layer-client-binary") + layer_stage = _layer_fault_stage(action["fault"]) if layer_seed_actions else None validate_fault_action( action, capabilities, f"{plan.path}:{action['id']}", capabilities_path, require_model=action["op"] == "crash", @@ -2294,6 +2456,7 @@ def run_daemon_fault_recovery( generation = 0 # next generation; fault, recovery, clean restart: 0, 1, 2 active_generation = 0 process: subprocess.Popen[str] | None = None + layer_client_process: subprocess.Popen[str] | None = None shm_names: list[str] = [] shm_base = f"/psharness_{os.getpid()}_{time.monotonic_ns()}" shm = "" @@ -2330,6 +2493,9 @@ def start_daemon(inject_fault: bool, action_id: str | None = None) -> subprocess "--nshards", str(plan.header["case"]["shards"]), "--storage", plan.header["case"]["storage"], ] + if layer_seed_actions: + command.extend(["--segment-size", "65536", "--flush-pages", "2", + "--compact-layers", "1000"]) env = private_environment() if inject_fault: # Keep these names local and explicit: inherited PAGESTORE_* values @@ -2402,9 +2568,22 @@ def wait_ready(daemon_process: subprocess.Popen[str]) -> dict[str, Any]: current_action_id = action["id"] emit("fault_arm", target="store", name=fault_name, hit=fault_hit) process = start_daemon(True, action["id"]) + if layer_seed_actions: + wait_ready(process) + layer_client_process = _start_layer_client( + layer_client, shm, "seed", trace / "layer-client.log" + ) deadline = time.monotonic() + fault_timeout if fault_action == "crash": while process.poll() is None and time.monotonic() < deadline: + if layer_client_process is not None and layer_client_process.poll() is not None: + if layer_client_process.returncode == 0: + raise FaultNotReached( + f"layer workload completed before fault {fault_name!r}" + ) + raise UnexpectedExit( + f"layer workload exited with status {layer_client_process.returncode}" + ) time.sleep(0.02) if process.poll() is None: raise FaultNotReached( @@ -2431,6 +2610,12 @@ def wait_ready(daemon_process: subprocess.Popen[str]) -> dict[str, Any]: scenario, seed, action["id"], ) shutil.copy2(report, trace / "fault-report.jsonl") + if layer_stage is not None: + _check_layer_crash_snapshot(store, layer_stage) + if layer_client_process is not None and layer_client_process.poll() is None: + layer_client_process.kill() + layer_client_process.wait(timeout=5) + layer_client_process = None emit("fault", target="store", name=fault_name, model=fault_model, returncode=fault_process.returncode, report=result, reached=True) remove_shm(shm) @@ -2551,14 +2736,29 @@ def wait_ready(daemon_process: subprocess.Popen[str]) -> dict[str, Any]: # Recovery is intentionally followed by one additional clean restart. process = start_daemon(False, action["id"]) health = wait_ready(process) + if layer_seed_actions: + _verify_layer_client(layer_client, shm, trace / "layer-client.log", timeout) + manifest = inspect_store(inspector, shm, "manifest", inspection_schema) + _check_layer_manifest(manifest, layer_stage, False) probe_runtime_inspection(inspector, shm, capabilities, inspection_schema) emit("recovered", target="store", health=health) stop_daemon() + if process.returncode != 0: + raise UnexpectedExit( + "recovered daemon did not stop cleanly before restart; status " + f"{process.returncode}" + ) emit("process_stop", target="store", pid=process.pid, returncode=process.returncode) remove_shm(shm) + process = None process = start_daemon(False, action["id"]) health = wait_ready(process) + if layer_seed_actions: + _verify_layer_client(layer_client, shm, trace / "layer-client.log", timeout) + _check_layer_manifest_after_restart( + inspector, shm, inspection_schema, layer_stage, timeout + ) probe_runtime_inspection(inspector, shm, capabilities, inspection_schema) emit("restarted", target="store", health=health) emit("run_pass") @@ -2593,6 +2793,8 @@ def wait_ready(daemon_process: subprocess.Popen[str]) -> dict[str, Any]: "--inspection-schema", str(bundle_inspection_schema), "--daemon-fault-recovery", str(root / "plan.jsonl"), "--daemon-binary", str(daemon), "--inspect-binary", str(inspector), + *( ["--layer-client-binary", str(layer_client)] + if layer_client is not None else [] ), "--run-root", str(root.parent / f"{root.name}.rerun"), "--keep", ], } @@ -2600,6 +2802,9 @@ def wait_ready(daemon_process: subprocess.Popen[str]) -> dict[str, Any]: json.dumps(metadata, indent=2, sort_keys=True) + "\n", encoding="utf-8" ) finally: + if layer_client_process is not None and layer_client_process.poll() is None: + layer_client_process.kill() + layer_client_process.wait(timeout=5) if process is not None: stop_daemon() emit("process_stop", target="store", pid=process.pid, @@ -4007,6 +4212,7 @@ def parse_args(argv: list[str]) -> argparse.Namespace: group.add_argument("--materializer-smoke", type=Path, metavar="PLAN") group.add_argument("--legacy-integration", action="store_true") parser.add_argument("--daemon-binary", type=Path) + parser.add_argument("--layer-client-binary", type=Path) parser.add_argument("--materializer-supervisor", type=Path) parser.add_argument("--build-dir", type=Path) parser.add_argument("--integration-script", type=Path, @@ -4134,6 +4340,7 @@ def main(argv: list[str] | None = None) -> int: ), args.daemon_binary, args.inspect_binary, args.run_root, args.keep, args.capabilities, command, + layer_client=args.layer_client_binary, ) if args.keep or args.run_root: print(root) diff --git a/contrib/pagestore/harness/scenarios/image_layer_after_create.jsonl b/contrib/pagestore/harness/scenarios/image_layer_after_create.jsonl new file mode 100644 index 0000000000000..00861208ccc27 --- /dev/null +++ b/contrib/pagestore/harness/scenarios/image_layer_after_create.jsonl @@ -0,0 +1,3 @@ +{"schema":1,"scenario":"image-layer-after-create","seed":2718281,"contracts":["fault_reachability","layer_sentinel","manifest_clean","restart_idempotence"],"case":{"storage":"posix","shards":1,"compute":["writer"]}} +{"op":"layer_seed","id":"seed-layer","target":"store"} +{"op":"crash","id":"image-layer-after-create","target":"store","model":"process_abort","fault":"image_layer.after_create","action":"crash","hit":1} diff --git a/contrib/pagestore/harness/scenarios/image_layer_after_manifest_add.jsonl b/contrib/pagestore/harness/scenarios/image_layer_after_manifest_add.jsonl new file mode 100644 index 0000000000000..a04b9e3968bcf --- /dev/null +++ b/contrib/pagestore/harness/scenarios/image_layer_after_manifest_add.jsonl @@ -0,0 +1,3 @@ +{"schema":1,"scenario":"image-layer-after-manifest-add","seed":2718284,"contracts":["fault_reachability","layer_sentinel","manifest_add","restart_idempotence"],"case":{"storage":"posix","shards":1,"compute":["writer"]}} +{"op":"layer_seed","id":"seed-layer","target":"store"} +{"op":"crash","id":"image-layer-after-manifest-add","target":"store","model":"process_abort","fault":"image_layer.after_manifest_add","action":"crash","hit":1} diff --git a/contrib/pagestore/harness/scenarios/image_layer_after_seal.jsonl b/contrib/pagestore/harness/scenarios/image_layer_after_seal.jsonl new file mode 100644 index 0000000000000..4c9966ebe3d2b --- /dev/null +++ b/contrib/pagestore/harness/scenarios/image_layer_after_seal.jsonl @@ -0,0 +1,3 @@ +{"schema":1,"scenario":"image-layer-after-seal","seed":2718283,"contracts":["fault_reachability","layer_sentinel","manifest_clean","restart_idempotence"],"case":{"storage":"posix","shards":1,"compute":["writer"]}} +{"op":"layer_seed","id":"seed-layer","target":"store"} +{"op":"crash","id":"image-layer-after-seal","target":"store","model":"process_abort","fault":"image_layer.after_seal","action":"crash","hit":1} diff --git a/contrib/pagestore/harness/scenarios/image_layer_after_write.jsonl b/contrib/pagestore/harness/scenarios/image_layer_after_write.jsonl new file mode 100644 index 0000000000000..bca0a72d5898c --- /dev/null +++ b/contrib/pagestore/harness/scenarios/image_layer_after_write.jsonl @@ -0,0 +1,3 @@ +{"schema":1,"scenario":"image-layer-after-write","seed":2718282,"contracts":["fault_reachability","layer_sentinel","manifest_clean","restart_idempotence"],"case":{"storage":"posix","shards":1,"compute":["writer"]}} +{"op":"layer_seed","id":"seed-layer","target":"store"} +{"op":"crash","id":"image-layer-after-write","target":"store","model":"process_abort","fault":"image_layer.after_write","action":"crash","hit":1} diff --git a/contrib/pagestore/harness/tests/test_plan.py b/contrib/pagestore/harness/tests/test_plan.py index f322d85c59520..8ebc953f174a9 100644 --- a/contrib/pagestore/harness/tests/test_plan.py +++ b/contrib/pagestore/harness/tests/test_plan.py @@ -523,6 +523,159 @@ def test_named_fault_catalog_validates_target_model_action_and_hit(self): MODULE.validate_plan(plan, capabilities) MODULE.validate_runtime_plan(plan, capabilities, "daemon_fault_smoke") + def test_image_layer_fault_scenarios_validate_as_ordered_composed_slices(self): + capabilities = MODULE.read_json(ROOT / "capabilities.json") + scenario_dir = ROOT / "scenarios" + scenarios = [ + scenario_dir / "image_layer_after_create.jsonl", + scenario_dir / "image_layer_after_write.jsonl", + scenario_dir / "image_layer_after_seal.jsonl", + scenario_dir / "image_layer_after_manifest_add.jsonl", + ] + for path in scenarios: + with self.subTest(scenario=path.name): + plan = MODULE.read_plan(path) + MODULE.validate_plan(plan, capabilities, ROOT / "capabilities.json") + MODULE.validate_runtime_plan(plan, capabilities, "daemon_fault_smoke") + + def test_image_layer_seed_cannot_follow_named_fault(self): + capabilities = MODULE.read_json(ROOT / "capabilities.json") + path = self.write_plan([ + { + "schema": 1, "scenario": "bad-layer-order", "seed": 1, + "contracts": ["fault_reachability"], + "case": {"storage": "posix", "shards": 1, "compute": ["writer"]}, + }, + { + "op": "crash", "id": "fault", "target": "store", + "model": "process_abort", "fault": "image_layer.after_create", + "action": "crash", "hit": 1, + }, + {"op": "layer_seed", "id": "seed", "target": "store"}, + ]) + plan = MODULE.read_plan(path) + MODULE.validate_plan(plan, capabilities, ROOT / "capabilities.json") + with self.assertRaisesRegex(MODULE.PlanError, "layer_seed before"): + MODULE.validate_runtime_plan(plan, capabilities, "daemon_fault_smoke") + + def test_manifest_add_recovery_allows_transient_compaction_states(self): + for layer_count in (1, 2): + with self.subTest(layer_count=layer_count): + MODULE._check_layer_manifest( + { + "layer_count": layer_count, + "local_layers": layer_count, + "manifest_poisoned": False, + }, + "manifest_add", + False, + ) + with self.assertRaisesRegex(MODULE.OracleMismatch, r"expected \(1, 1\)"): + MODULE._check_layer_manifest( + {"layer_count": 2, "local_layers": 2, "manifest_poisoned": False}, + "manifest_add", + True, + ) + with self.assertRaisesRegex( + MODULE.OracleMismatch, r"expected \(1, 1\) or \(2, 2\)" + ): + MODULE._check_layer_manifest( + {"layer_count": 0, "local_layers": 0, "manifest_poisoned": False}, + "manifest_add", + False, + ) + + def test_manifest_add_restart_polls_until_compaction_converges(self): + manifests = [ + {"layer_count": 2, "local_layers": 2, "manifest_poisoned": False}, + {"layer_count": 1, "local_layers": 1, "manifest_poisoned": False}, + ] + with ( + mock.patch.object(MODULE, "inspect_store", side_effect=manifests) as inspect, + mock.patch.object(MODULE.time, "monotonic", side_effect=[10.0, 10.01]), + mock.patch.object(MODULE.time, "sleep") as sleep, + ): + MODULE._check_layer_manifest_after_restart( + Path("inspect"), "/shm", {}, "manifest_add", 1.0 + ) + self.assertEqual(inspect.call_count, 2) + sleep.assert_called_once_with(0.05) + + def test_manifest_add_restart_accepts_immediate_convergence(self): + with ( + mock.patch.object( + MODULE, + "inspect_store", + return_value={ + "layer_count": 1, + "local_layers": 1, + "manifest_poisoned": False, + }, + ) as inspect, + mock.patch.object(MODULE.time, "monotonic", return_value=10.0), + mock.patch.object(MODULE.time, "sleep") as sleep, + ): + MODULE._check_layer_manifest_after_restart( + Path("inspect"), "/shm", {}, "manifest_add", 1.0 + ) + inspect.assert_called_once() + sleep.assert_not_called() + + def test_manifest_add_restart_times_out_on_persistent_duplicate(self): + with ( + mock.patch.object( + MODULE, + "inspect_store", + return_value={ + "layer_count": 2, + "local_layers": 2, + "manifest_poisoned": False, + }, + ) as inspect, + mock.patch.object(MODULE.time, "monotonic", side_effect=[10.0, 10.1]), + mock.patch.object(MODULE.time, "sleep") as sleep, + ): + with self.assertRaisesRegex(MODULE.HarnessTimeout, "did not converge"): + MODULE._check_layer_manifest_after_restart( + Path("inspect"), "/shm", {}, "manifest_add", 0.05 + ) + inspect.assert_called_once() + sleep.assert_not_called() + + def test_manifest_add_restart_rejects_invalid_or_poisoned_state_immediately(self): + cases = [ + ( + {"layer_count": 2, "local_layers": 1, "manifest_poisoned": False}, + r"expected \(1, 1\)", + ), + ( + {"layer_count": 2, "local_layers": 2, "manifest_poisoned": True}, + "manifest_poisoned=True", + ), + ] + for manifest, message in cases: + with self.subTest(manifest=manifest): + with mock.patch.object(MODULE, "inspect_store", return_value=manifest): + with self.assertRaisesRegex(MODULE.OracleMismatch, message): + MODULE._check_layer_manifest_after_restart( + Path("inspect"), "/shm", {}, "manifest_add", 1.0 + ) + + def test_layer_client_uses_absolute_executable_path(self): + with tempfile.TemporaryDirectory() as temporary: + root = Path(temporary) + client = Path("pagestore_layer_crash_client") + expected = str(client.resolve()) + result = mock.Mock(returncode=0) + with ( + mock.patch.object(MODULE.subprocess, "Popen") as popen, + mock.patch.object(MODULE.subprocess, "run", return_value=result) as run, + ): + MODULE._start_layer_client(client, "/shm", "seed", root / "seed.log") + MODULE._verify_layer_client(client, "/shm", root / "verify.log", 1.0) + self.assertEqual(popen.call_args.args[0][0], expected) + self.assertEqual(run.call_args.args[0][0], expected) + def test_named_fault_catalog_rejects_wrong_hit_without_launch(self): capabilities = MODULE.read_json(ROOT / "capabilities.json") path = self.write_plan([ @@ -907,7 +1060,7 @@ def fake_popen(command, **kwargs): MODULE.run_daemon_fault_recovery( plan, capabilities, inspection_schema, Path("daemon"), Path("inspect"), root, True, - ROOT / "capabilities.json", + ROOT / "capabilities.json", layer_client=Path("layer-client"), ) metadata = json.loads((root / "failure.json").read_text(encoding="utf-8")) @@ -922,8 +1075,90 @@ def fake_popen(command, **kwargs): (root / "fault-control" / "report.jsonl").read_text(encoding="utf-8"), "{malformed\n", ) + layer_client_index = metadata["command"].index("--layer-client-binary") + self.assertEqual( + metadata["command"][layer_client_index + 1], + str(Path("layer-client").resolve()), + ) self.assertFalse((root / "fault-control" / "arm").exists()) + def test_layer_recovery_requires_clean_shutdown_before_restart(self): + capabilities = MODULE.read_json(ROOT / "capabilities.json") + plan = MODULE.read_plan(ROOT / "scenarios" / "image_layer_after_create.jsonl") + directory = tempfile.TemporaryDirectory() + self.addCleanup(directory.cleanup) + root = Path(directory.name) / "failure" + health = { + "protocol_version": 45, "page_size": 8192, "io_unit": 262144, + "nchannels": 128, "nshards": 1, "admission_fence_epoch": 0, + "admission_pending_epoch": 0, "admission_pending_lsn": 0, + } + processes = [] + + class FakeProcess: + next_pid = 800 + + def __init__(self, env, shutdown_status=0): + self.pid = FakeProcess.next_pid + FakeProcess.next_pid += 1 + self.returncode = None + self.shutdown_status = shutdown_status + + def poll(self): + return self.returncode + + def wait(self, timeout=None): + if self.returncode is None: + self.returncode = self.shutdown_status + return self.returncode + + def kill(self): + self.returncode = -9 + + fault_control = None + + def fake_popen(command, **kwargs): + nonlocal fault_control + process = FakeProcess( + kwargs["env"], shutdown_status=1 if len(processes) == 1 else 0 + ) + processes.append(process) + if "PAGESTORE_TEST_FAULT_NAME" in kwargs["env"]: + fault_control = Path(kwargs["env"]["PAGESTORE_TEST_FAULT_DIR"]) + return process + + def trigger_fault(*args, **kwargs): + processes[0].returncode = 88 + assert fault_control is not None + (fault_control / "report.jsonl").write_text( + json.dumps({ + "schema": 1, "name": "image_layer.after_create", + "action": "crash", "hit": 1, "pid": processes[0].pid, + }) + "\n", + encoding="utf-8", + ) + + with ( + mock.patch.object(MODULE.subprocess, "Popen", side_effect=fake_popen), + mock.patch.object(MODULE, "_start_layer_client", side_effect=trigger_fault), + mock.patch.object(MODULE, "inspect_store", return_value=health), + mock.patch.object(MODULE, "validate_runtime_health"), + mock.patch.object(MODULE, "probe_runtime_inspection"), + mock.patch.object(MODULE, "_check_layer_crash_snapshot"), + mock.patch.object(MODULE, "_verify_layer_client"), + mock.patch.object(MODULE, "_check_layer_manifest"), + mock.patch.object(MODULE, "signal_process_group"), + mock.patch.object(MODULE, "remove_shm"), + ): + with self.assertRaisesRegex( + MODULE.PlanError, "did not stop cleanly before restart; status 1" + ): + MODULE.run_daemon_fault_recovery( + plan, capabilities, {}, Path("daemon"), Path("inspect"), + root, True, ROOT / "capabilities.json", layer_client=Path("client"), + ) + self.assertEqual(len(processes), 2) + def test_fault_report_rejects_extra_fields_and_multiple_lines(self): with tempfile.TemporaryDirectory() as temporary: report = Path(temporary) / "report.jsonl" diff --git a/contrib/pagestore/meson.build b/contrib/pagestore/meson.build index fef119974db65..0b728f0f35028 100644 --- a/contrib/pagestore/meson.build +++ b/contrib/pagestore/meson.build @@ -81,6 +81,14 @@ pagestore_daemon = executable('pagestore_daemon', ) contrib_targets += pagestore_daemon +# Deterministic POSIX image-layer crash workload/oracle used by the H1 harness +# slice. It only speaks the standalone daemon IPC protocol. +pagestore_layer_crash_client = executable('pagestore_layer_crash_client', + files('pagestore_layer_crash_client.c'), + kwargs: default_bin_args, +) +contrib_targets += pagestore_layer_crash_client + # Pure process-local unit test for the first crash-only named fault registry. pagestore_fault_test = executable('pagestore_fault_test', files('pagestore_fault_test.c', 'pagestore_fault.c'), @@ -176,7 +184,9 @@ test('pagestore_inspect_mailbox', # Unit test for the immutable image-layer file format (no daemon, no PG). pagestore_layer_test = executable('pagestore_layer_test', - files('pagestore_layer_test.c', 'pagestore_layer.c', 'pagestore_layer_store.c'), + files('pagestore_layer_test.c', 'pagestore_layer.c', 'pagestore_layer_store.c', + 'pagestore_fault.c'), + dependencies: [dependency('threads')], install: false, ) test('pagestore_layer', @@ -187,7 +197,9 @@ test('pagestore_layer', # Unit test for the filesystem-backed remote object provider (no daemon, no PG). pagestore_layer_store_test = executable('pagestore_layer_store_test', - files('pagestore_layer_store_test.c', 'pagestore_layer_store.c', 'pagestore_layer.c'), + files('pagestore_layer_store_test.c', 'pagestore_layer_store.c', 'pagestore_layer.c', + 'pagestore_fault.c'), + dependencies: [dependency('threads')], install: false, ) test('pagestore_layer_store', @@ -334,7 +346,8 @@ test('pagestore_backpressure_daemon', # Unit test for the durable layer manifest replay (no daemon, no PG). pagestore_manifest_test = executable('pagestore_manifest_test', files('pagestore_manifest_test.c', 'pagestore_manifest.c', 'pagestore_layer.c', - 'pagestore_layer_store.c'), + 'pagestore_layer_store.c', 'pagestore_fault.c'), + dependencies: [dependency('threads')], install: false, ) test('pagestore_manifest', @@ -358,7 +371,8 @@ test('pagestore_retention', # Unit test for crash-safe compaction/GC recovery (no daemon, no PG). pagestore_gc_test = executable('pagestore_gc_test', files('pagestore_gc_test.c', 'pagestore_manifest.c', 'pagestore_layer.c', - 'pagestore_layer_store.c'), + 'pagestore_layer_store.c', 'pagestore_fault.c'), + dependencies: [dependency('threads')], install: false, ) test('pagestore_gc', @@ -577,3 +591,22 @@ foreach fault_scenario : ['daemon_fault_error', 'daemon_fault_pause'] timeout: 30, ) endforeach +foreach layer_fault_scenario : [ + 'image_layer_after_create', + 'image_layer_after_write', + 'image_layer_after_seal', + 'image_layer_after_manifest_add', +] + test('pagestore_harness_' + layer_fault_scenario, + pagestore_harness_python, + args: [files('harness/pagestore_harness.py'), + '--capabilities', files('harness/capabilities.json'), + '--daemon-fault-recovery', files('harness/scenarios/' + layer_fault_scenario + '.jsonl'), + '--daemon-binary', pagestore_daemon, + '--layer-client-binary', pagestore_layer_crash_client, + '--inspect-binary', pagestore_inspect], + depends: [pagestore_daemon, pagestore_layer_crash_client, pagestore_inspect], + suite: ['pagestore'], + timeout: 60, + ) +endforeach diff --git a/contrib/pagestore/pagestore_fault_points.def b/contrib/pagestore/pagestore_fault_points.def index c3a1d1637bda7..1478b4bd0a185 100644 --- a/contrib/pagestore/pagestore_fault_points.def +++ b/contrib/pagestore/pagestore_fault_points.def @@ -20,3 +20,10 @@ PAGESTORE_FAULT_POINT(BRANCH_PREPARE_BEFORE_PREPARED_RECEIPT, "branch_prepare.be PAGESTORE_FAULT_POINT(BRANCH_PREPARE_AFTER_PREPARED_RECEIPT, "branch_prepare.after_prepared_receipt", "branch", "process_abort", "crash", 1, 1) PAGESTORE_FAULT_POINT(BRANCH_PREPARE_AFTER_MATERIALIZER_RESUME, "branch_prepare.after_materializer_resume", "branch", "process_abort", "crash", 1, 1) PAGESTORE_FAULT_POINT(BRANCH_PREPARE_AFTER_WRITER_RESTORE, "branch_prepare.after_writer_restore", "branch", "process_abort", "crash", 1, 1) +// Image-layer H1 probes are ordered POSIX-local process-abort boundaries: the +// file exists, has been written but not sealed, has been fsync'd, and its ADD +// has been fsync'd in layers.manifest. after_write is not a power-loss claim. +PAGESTORE_FAULT_POINT(IMAGE_LAYER_AFTER_CREATE, "image_layer.after_create", "store", "process_abort", "crash", 1, 1) +PAGESTORE_FAULT_POINT(IMAGE_LAYER_AFTER_WRITE, "image_layer.after_write", "store", "process_abort", "crash", 1, 1) +PAGESTORE_FAULT_POINT(IMAGE_LAYER_AFTER_SEAL, "image_layer.after_seal", "store", "process_abort", "crash", 1, 1) +PAGESTORE_FAULT_POINT(IMAGE_LAYER_AFTER_MANIFEST_ADD, "image_layer.after_manifest_add", "store", "process_abort", "crash", 1, 1) diff --git a/contrib/pagestore/pagestore_fault_test.c b/contrib/pagestore/pagestore_fault_test.c index e5ae2df805239..4b5b83026b5fd 100644 --- a/contrib/pagestore/pagestore_fault_test.c +++ b/contrib/pagestore/pagestore_fault_test.c @@ -158,7 +158,13 @@ main(void) point == PS_FAULT_POINT_DAEMON_AFTER_READY && strcmp(ps_fault_allowed_actions(point), "crash|error|pause") == 0 && strcmp(ps_fault_allowed_actions(PS_FAULT_POINT_PAGE_PRUNE_AFTER_FRONTIER), - "crash") == 0, "catalog exposes action policy"); + "crash") == 0 && + ps_fault_lookup("image_layer.after_create", &point) == 0 && + strcmp(ps_fault_allowed_actions(point), "crash") == 0 && + ps_fault_lookup("image_layer.after_write", &point) == 0 && + ps_fault_lookup("image_layer.after_seal", &point) == 0 && + ps_fault_lookup("image_layer.after_manifest_add", &point) == 0, + "catalog exposes action policy and H1 image-layer points"); /* Partial configuration, bad hit values, and non-crash lock-held actions fail closed. */ setenv("PAGESTORE_TEST_FAULT_ACTION", "crash", 1); diff --git a/contrib/pagestore/pagestore_layer_crash_client.c b/contrib/pagestore/pagestore_layer_crash_client.c new file mode 100644 index 0000000000000..96f158a61e4d4 --- /dev/null +++ b/contrib/pagestore/pagestore_layer_crash_client.c @@ -0,0 +1,210 @@ +/*------------------------------------------------------------------------- + * + * pagestore_layer_crash_client.c + * Deterministic IPC workload/oracle for the POSIX image-layer H1 crash slice. + * + * The seed writes two relation pages. With --flush-pages 2 the second write + * drives exactly one image-layer publication. The verify mode checks both + * sentinel bytes and the resolved page LSN after recovery. + * + *------------------------------------------------------------------------- + */ +#include +#include +#include +#include +#include +#include +#include + +#include "pagestore_ipc.h" + +#define TEST_REL 4242u +#define TEST_LSN0 UINT64_C(0x1000) +#define TEST_LSN1 UINT64_C(0x2000) + +static void *shm_base; +static int shm_fd = -1; +static int channel = -1; +static uint32_t page_size; + +static void +die(const char *message) +{ + fprintf(stderr, "pagestore_layer_crash_client: %s\n", message); + exit(1); +} + +static void +attach(const char *name) +{ + PsShmHeader *header; + + shm_fd = shm_open(name, O_RDWR, 0600); + if (shm_fd < 0) + die("cannot open shared memory"); + shm_base = mmap(NULL, PS_SHM_SIZE, PROT_READ | PROT_WRITE, MAP_SHARED, + shm_fd, 0); + if (shm_base == MAP_FAILED) + die("cannot map shared memory"); + header = (PsShmHeader *) shm_base; + if (header->magic != PS_SHM_MAGIC || + ps_load_acquire(&header->startup_state) != PS_SHM_READY || + header->nshards != 1) + die("daemon is not ready for the single-shard H1 workload"); + page_size = header->page_size; + if (page_size == 0 || page_size > PS_IO_UNIT) + die("invalid daemon page size"); + for (uint32_t i = 0; i < header->nchannels; i++) + if (ps_cas(&ps_channel(shm_base, i)->claimed, 0, 1)) + { + channel = (int) i; + return; + } + die("no free daemon channel"); +} + +static void +detach(void) +{ + if (shm_base != NULL && shm_base != MAP_FAILED) + { + if (channel >= 0) + ps_store_release(&ps_channel(shm_base, channel)->claimed, 0); + munmap(shm_base, PS_SHM_SIZE); + } + if (shm_fd >= 0) + close(shm_fd); +} + +static PsChannel * +execute(void) +{ + PsChannel *ch = ps_channel(shm_base, channel); + + ps_request_generation_next(ch); + ps_store_release(&ch->state, PS_STATE_REQUEST); + while (ps_load_acquire(&ch->state) != PS_STATE_DONE) + ; + return ch; +} + +static void +set_relation(PsChannel *ch) +{ + memset(&ch->key, 0, sizeof(ch->key)); + ch->key.spcOid = 1; + ch->key.dbOid = 1; + ch->key.relNumber = TEST_REL; + ch->key.forkNum = 0; + ch->key.klass = PS_KLASS_RELATION; + ch->timeline = 0; + ch->req_lsn = 0; + ch->req_seq = 0; + ch->incarnation = 0; +} + +static void +fill_page(unsigned char *page, uint64_t lsn, unsigned char tag) +{ + uint32_t high = (uint32_t) (lsn >> 32); + uint32_t low = (uint32_t) lsn; + + memcpy(page, &high, sizeof(high)); + memcpy(page + sizeof(high), &low, sizeof(low)); + for (uint32_t i = 8; i < page_size; i++) + page[i] = (unsigned char) (tag ^ (i & 0xff)); +} + +static int +page_matches(const unsigned char *page, uint64_t lsn, unsigned char tag) +{ + uint32_t high; + uint32_t low; + + memcpy(&high, page, sizeof(high)); + memcpy(&low, page + sizeof(high), sizeof(low)); + if ((((uint64_t) high << 32) | low) != lsn) + return 0; + for (uint32_t i = 8; i < page_size; i++) + if (page[i] != (unsigned char) (tag ^ (i & 0xff))) + return 0; + return 1; +} + +static void +seed(void) +{ + PsChannel *ch = ps_channel(shm_base, channel); + unsigned char *pages = malloc((size_t) page_size * 2); + + if (pages == NULL) + die("out of memory"); + fill_page(pages, TEST_LSN0, 0xa1); + fill_page(pages + page_size, TEST_LSN1, 0xb2); + set_relation(ch); + ch->opcode = PS_OP_CREATE; + if (execute()->status != PS_STATUS_OK) + die("relation create failed"); + set_relation(ch); + ch->opcode = PS_OP_WRITEV; + ch->blocknum = 0; + ch->nblocks = 2; + memcpy(ch->data, pages, (size_t) page_size * 2); + if (execute()->status != PS_STATUS_OK) + die("layer-seeding write failed"); + free(pages); +} + +static void +verify(void) +{ + PsChannel *ch = ps_channel(shm_base, channel); + unsigned char *page = malloc(page_size); + const uint64_t lsns[] = {TEST_LSN0, TEST_LSN1}; + const unsigned char tags[] = {0xa1, 0xb2}; + + if (page == NULL) + die("out of memory"); + for (uint32_t block = 0; block < 2; block++) + { + set_relation(ch); + ch->opcode = PS_OP_READ_AT; + ch->blocknum = block; + ch->req_lsn = TEST_LSN1 + 1; + if (execute()->status != PS_STATUS_OK || ch->result != 1 || + ch->req_lsn != lsns[block]) + die("sentinel read did not return the expected page LSN"); + memcpy(page, ch->data, page_size); + if (!page_matches(page, lsns[block], tags[block])) + die("sentinel page contents changed across recovery"); + } + free(page); +} + +int +main(int argc, char **argv) +{ + const char *shm = NULL; + const char *mode = NULL; + + for (int i = 1; i < argc; i++) + { + if (strcmp(argv[i], "--shm") == 0 && i + 1 < argc) + shm = argv[++i]; + else if (strcmp(argv[i], "--mode") == 0 && i + 1 < argc) + mode = argv[++i]; + else + die("usage: --shm NAME --mode seed|verify"); + } + if (shm == NULL || mode == NULL || + (strcmp(mode, "seed") != 0 && strcmp(mode, "verify") != 0)) + die("usage: --shm NAME --mode seed|verify"); + attach(shm); + if (strcmp(mode, "seed") == 0) + seed(); + else + verify(); + detach(); + return 0; +} diff --git a/contrib/pagestore/pagestore_layer_store.c b/contrib/pagestore/pagestore_layer_store.c index d7b4fc3ceb007..eafee14846dbf 100644 --- a/contrib/pagestore/pagestore_layer_store.c +++ b/contrib/pagestore/pagestore_layer_store.c @@ -19,6 +19,7 @@ #include #include +#include "pagestore_fault.h" #include "pagestore_layer_store.h" static uint32_t layer_page_size = PS_DEFAULT_PAGE_SIZE; @@ -448,6 +449,8 @@ local_create_local_layer(uint64_t layer_id, char *uri, uint32_t uri_len) unlink(path); return -1; } + if (ps_fault_probe(PS_FAULT_POINT_IMAGE_LAYER_AFTER_CREATE) != 0) + return -1; n = snprintf(uri, uri_len, "%s", path); if (n < 0 || (uint32_t) n >= uri_len) { @@ -494,7 +497,7 @@ local_write_local_layer(uint64_t layer_id, const void *buf, uint64_t len) done += (uint64_t) w; } close(fd); - return 0; + return ps_fault_probe(PS_FAULT_POINT_IMAGE_LAYER_AFTER_WRITE) == 0 ? 0 : -1; } static int @@ -511,7 +514,9 @@ local_seal_local_layer(uint64_t layer_id) return -1; rc = fsync(fd); close(fd); - return rc; + if (rc != 0) + return rc; + return ps_fault_probe(PS_FAULT_POINT_IMAGE_LAYER_AFTER_SEAL) == 0 ? 0 : -1; } static int diff --git a/contrib/pagestore/pagestore_manifest.c b/contrib/pagestore/pagestore_manifest.c index 8651d261c16fb..48cf3989b0ea0 100644 --- a/contrib/pagestore/pagestore_manifest.c +++ b/contrib/pagestore/pagestore_manifest.c @@ -18,6 +18,7 @@ #include #include "pagestore_manifest.h" +#include "pagestore_fault.h" #define PS_MANIFEST_MAGIC 0x504d414e /* "PMAN" */ #define PS_MANIFEST_VERSION 3 /* 3: per-record CRC (over the header + payload) */ @@ -793,6 +794,10 @@ ps_manifest_add_layer(const PsLayerDesc *desc) return -1; if (manifest_append(PS_MANIFEST_ADD_LAYER, &disk, sizeof(disk)) != 0) return -1; + /* The ADD record is now durable. The in-memory map update follows this + * probe, so recovery must be able to choose the replayed old/new state. */ + if (ps_fault_probe(PS_FAULT_POINT_IMAGE_LAYER_AFTER_MANIFEST_ADD) != 0) + return -1; return ps_layer_map_add(&ps_layer_map, desc); }