From 38739fdfe528a299b5e5617c0b957882ef517f67 Mon Sep 17 00:00:00 2001 From: venkat1701 Date: Mon, 5 Oct 2026 05:17:01 +0530 Subject: [PATCH 1/3] feat(bench): time recovery after killing an ingest process A child JVM ingests nodes and prints every acknowledged commit. The parent kills it with SIGKILL halfway, times opening the directory and counts how many acknowledged nodes are missing. HStoreStore can now open an existing directory without redefining its types. --- .../main/java/io/hstore/bench/Comparison.java | 75 +++++++++++++++++-- .../java/io/hstore/bench/HStoreStore.java | 23 ++++-- .../io/hstore/bench/HyperGraphDbStore.java | 9 ++- .../src/main/java/io/hstore/bench/Store.java | 10 ++- 4 files changed, 102 insertions(+), 15 deletions(-) diff --git a/benchmarks/src/main/java/io/hstore/bench/Comparison.java b/benchmarks/src/main/java/io/hstore/bench/Comparison.java index 7998869..07593ad 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Comparison.java +++ b/benchmarks/src/main/java/io/hstore/bench/Comparison.java @@ -2,8 +2,10 @@ import io.hstore.db.value.Json; +import java.io.BufferedReader; import java.io.IOException; import java.io.UncheckedIOException; +import java.lang.management.ManagementFactory; import java.math.BigDecimal; import java.nio.file.Files; import java.nio.file.Path; @@ -63,6 +65,7 @@ private interface Slice { private static final int BATCH = 1_000; private static final int COMMITS = 5_000; + private static final String COMMITTED = "committed "; private static final long MIXED_NANOS = 10_000_000_000L; private static final int MIXED_BATCH = 10; private static final int MIXED_UPDATES_PER_SECOND = 1_000; @@ -76,10 +79,12 @@ public static void main(String[] args) throws Exception { Map options = options(args); switch (args.length == 0 ? "" : args[0]) { case "run" -> run(options); + case "crash" -> crash(options); case "report" -> IO.println(Report.markdown(args)); default -> { IO.println(""" usage: Comparison run --store hstore|hypergraphdb [--scale N] [--durability async|sync] [--threads N] [--history N] [--cache-mb N] --out FILE + Comparison crash --store hstore|hypergraphdb --dir DIRECTORY (used by run; ingests until killed) Comparison report FILE... (runs of two stores; cells show the median and range)"""); System.exit(2); } @@ -103,16 +108,11 @@ private static void run(Map options) throws Exception { int threads = Integer.parseInt(options.getOrDefault("threads", String.valueOf(Runtime.getRuntime().availableProcessors()))); Path out = Path.of(options.getOrDefault("out", "results/" + storeName + ".json")); Path directory = Files.createTempDirectory("hstore-bench-" + storeName); - long cacheBytes = Long.parseLong(options.getOrDefault("cache-mb", "0")) << 20; Dataset dataset = Dataset.generate(scale, SEED); List results = new ArrayList<>(); Map totals = new LinkedHashMap<>(); long baselineHeap = JvmUsage.retainedHeap(); - Store store = switch (storeName) { - case "hstore" -> new HStoreStore(directory, sync, Integer.parseInt(options.getOrDefault("history", "64")), cacheBytes); - case "hypergraphdb" -> new HyperGraphDbStore(directory, sync, cacheBytes); - default -> throw new IllegalArgumentException("unknown store " + storeName); - }; + Store store = open(options, directory, true); IO.println("%s: %,d nodes, %,d hyperedges, %,d incidences, durability %s, %d threads" .formatted(store.name(), dataset.nodes(), dataset.edges().length, dataset.incidences(), sync ? "sync" : "async", threads)); try { @@ -196,11 +196,74 @@ private static void run(Map options) throws Exception { } long churned = store.diskBytes(); results.add(new Result("disk.churned", "bytes on disk after the deletes and a clean shutdown", churned, 1000.0, churned)); + results.add(recover(options, dataset)); write(out, store, scale, sync, threads, dataset, totals, results); results.forEach(result -> IO.println(" %-22s %,14.0f ops/s %,10.1f ms checksum %d" .formatted(result.workload(), result.rate(), result.millis(), result.checksum()))); } + private static Store open(Map options, Path directory, boolean create) { + boolean sync = options.getOrDefault("durability", "async").equals("sync"); + long cacheBytes = Long.parseLong(options.getOrDefault("cache-mb", "0")) << 20; + return switch (options.getOrDefault("store", "hstore")) { + case "hstore" -> new HStoreStore(directory, sync, Integer.parseInt(options.getOrDefault("history", "64")), cacheBytes, create); + case "hypergraphdb" -> new HyperGraphDbStore(directory, sync, cacheBytes); + case String other -> throw new IllegalArgumentException("unknown store " + other); + }; + } + + private static void crash(Map options) throws InterruptedException { + Dataset dataset = Dataset.generate(Integer.parseInt(options.getOrDefault("scale", "1")), SEED); + Store store = open(options, Path.of(options.get("dir")), true); + store.ingestNodes(dataset, BATCH, committed -> { + System.out.println(COMMITTED + committed); + System.out.flush(); + }); + Thread.sleep(Long.MAX_VALUE); + } + + private static Result recover(Map options, Dataset dataset) throws Exception { + Path directory = Files.createTempDirectory("hstore-bench-crash"); + List command = new ArrayList<>(); + command.add(ProcessHandle.current().info().command().orElse("java")); + command.addAll(ManagementFactory.getRuntimeMXBean().getInputArguments()); + command.addAll(List.of("-cp", System.getProperty("java.class.path"), Comparison.class.getName(), "crash", "--dir", directory.toString())); + for (String option : List.of("store", "scale", "durability", "history", "cache-mb")) { + if (options.containsKey(option)) { + command.addAll(List.of("--" + option, options.get(option))); + } + } + Process child = new ProcessBuilder(command).redirectErrorStream(true).start(); + int target = dataset.nodes() / 2; + int acknowledged = 0; + try (BufferedReader lines = child.inputReader()) { + for (String line = lines.readLine(); line != null && acknowledged < target; line = lines.readLine()) { + if (line.startsWith(COMMITTED)) { + acknowledged = Integer.parseInt(line.substring(COMMITTED.length())); + } + } + } finally { + child.destroyForcibly(); + child.waitFor(); + } + if (acknowledged < target) { + throw new IllegalStateException("the crash child exited after %,d acknowledged nodes".formatted(acknowledged)); + } + Store[] recovered = new Store[1]; + Result result = measure("recover", "open the store after killing the ingest process with SIGKILL", 1, () -> { + recovered[0] = open(options, directory, false); + return 1; + }); + long present; + try { + present = recovered[0].countNodes(); + } finally { + recovered[0].close(); + } + return result.with("checked", new Json.Bool(false)).with("acknowledged", number(acknowledged)) + .with("lost", number(Math.max(0, acknowledged - present))); + } + private static Json number(long value) { return new Json.Number(BigDecimal.valueOf(value)); } diff --git a/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java b/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java index 236d666..66c70f5 100644 --- a/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java +++ b/benchmarks/src/main/java/io/hstore/bench/HStoreStore.java @@ -16,6 +16,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.function.IntConsumer; final class HStoreStore implements Store { @@ -37,7 +38,7 @@ final class HStoreStore implements Store { private final int cachedNodes; private final long cacheBytes; - HStoreStore(Path directory, boolean sync, int history, long cacheBytes) { + HStoreStore(Path directory, boolean sync, int history, long cacheBytes, boolean create) { this.directory = directory; this.history = history; this.cacheBytes = cacheBytes; @@ -46,11 +47,13 @@ final class HStoreStore implements Store { this.options = DatabaseOptions.defaults().withEngine(engine .withDurability(sync ? Durability.SYNC : Durability.ASYNC).withCachedNodes(cachedNodes).withHistoryLimit(history)); this.database = HypergraphDatabase.open(directory, options); - database.write(writer -> { - writer.defineNode(NODE, List.of(new PropertyDef("name", TypeTag.STRING, false, false))); - writer.defineEdge(EDGE, AtomKind.SET_EDGE, List.of(), List.of()); - return null; - }); + if (create) { + database.write(writer -> { + writer.defineNode(NODE, List.of(new PropertyDef("name", TypeTag.STRING, false, false))); + writer.defineEdge(EDGE, AtomKind.SET_EDGE, List.of(), List.of()); + return null; + }); + } } @Override @@ -75,7 +78,7 @@ public Path directory() { } @Override - public void ingestNodes(Dataset dataset, int batch) { + public void ingestNodes(Dataset dataset, int batch, IntConsumer committed) { nodes = new long[dataset.nodes()]; for (int start = 0; start < nodes.length; start += batch) { int from = start; @@ -86,9 +89,15 @@ public void ingestNodes(Dataset dataset, int batch) { } return null; }); + committed.accept(to); } } + @Override + public long countNodes() { + return database.read(reader -> reader.atoms(NODE).count()); + } + @Override public void ingestEdges(Dataset dataset, int batch) { int[][] source = dataset.edges(); diff --git a/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java b/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java index 4fec7a9..5d48958 100644 --- a/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java +++ b/benchmarks/src/main/java/io/hstore/bench/HyperGraphDbStore.java @@ -21,6 +21,7 @@ import java.util.HashSet; import java.util.Set; import java.util.concurrent.Callable; +import java.util.function.IntConsumer; final class HyperGraphDbStore implements Store { @@ -78,7 +79,7 @@ private T read(Callable work) { } @Override - public void ingestNodes(Dataset dataset, int batch) { + public void ingestNodes(Dataset dataset, int batch, IntConsumer committed) { nodes = new HGPersistentHandle[dataset.nodes()]; for (int start = 0; start < nodes.length; start += batch) { int from = start; @@ -89,9 +90,15 @@ public void ingestNodes(Dataset dataset, int batch) { } return null; }); + committed.accept(to); } } + @Override + public long countNodes() { + return read(() -> hg.count(graph, hg.type(String.class))); + } + @Override public void ingestEdges(Dataset dataset, int batch) { int[][] source = dataset.edges(); diff --git a/benchmarks/src/main/java/io/hstore/bench/Store.java b/benchmarks/src/main/java/io/hstore/bench/Store.java index 01898aa..8287e8d 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Store.java +++ b/benchmarks/src/main/java/io/hstore/bench/Store.java @@ -4,6 +4,7 @@ import java.io.UncheckedIOException; import java.nio.file.Files; import java.nio.file.Path; +import java.util.function.IntConsumer; import java.util.stream.Stream; interface Store extends AutoCloseable { @@ -14,7 +15,14 @@ interface Store extends AutoCloseable { String cache(); - void ingestNodes(Dataset dataset, int batch); + default void ingestNodes(Dataset dataset, int batch) { + ingestNodes(dataset, batch, _ -> { + }); + } + + void ingestNodes(Dataset dataset, int batch, IntConsumer committed); + + long countNodes(); void ingestEdges(Dataset dataset, int batch); From 114b1db2cd9ae393ec4904508431d2ab5750cc49 Mon Sep 17 00:00:00 2001 From: venkat1701 Date: Mon, 5 Oct 2026 05:17:01 +0530 Subject: [PATCH 2/3] feat(bench): show lost acknowledged nodes even when there are none --- benchmarks/src/main/java/io/hstore/bench/Report.java | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/benchmarks/src/main/java/io/hstore/bench/Report.java b/benchmarks/src/main/java/io/hstore/bench/Report.java index 3f2d50d..a81eaab 100644 --- a/benchmarks/src/main/java/io/hstore/bench/Report.java +++ b/benchmarks/src/main/java/io/hstore/bench/Report.java @@ -71,12 +71,13 @@ static String markdown(String[] args) { return out.toString(); } - private record Metric(String key, String name, boolean perOperation, DoubleFunction format) { + private record Metric(String key, String name, boolean perOperation, boolean showZero, DoubleFunction format) { } - private static final List METRICS = List.of(new Metric("bytesWritten", "bytes written per operation", true, Report::bytes), - new Metric("allocated", "heap allocated per operation", true, Report::bytes), - new Metric("gcMillis", "GC pause time", false, value -> "%,.0f ms".formatted(value))); + private static final List METRICS = List.of(new Metric("bytesWritten", "bytes written per operation", true, false, Report::bytes), + new Metric("allocated", "heap allocated per operation", true, false, Report::bytes), + new Metric("gcMillis", "GC pause time", false, false, value -> "%,.0f ms".formatted(value)), + new Metric("lost", "acknowledged nodes lost", false, true, value -> "%,.0f".formatted(value))); private static void resources(StringBuilder out, List subjects, List baselines) { List all = Stream.concat(subjects.stream(), baselines.stream()).toList(); @@ -86,7 +87,7 @@ private static void resources(StringBuilder out, List subjects, List b if (all.stream().allMatch(run -> run.results().containsKey(workload) && run.results().get(workload).containsKey(metric.key()))) { double mine = median(metric(subjects, workload, metric)); double theirs = median(metric(baselines, workload, metric)); - if (mine == 0 && theirs == 0) { + if (mine == 0 && theirs == 0 && !metric.showZero()) { continue; } double ratio = theirs / mine; From e3d490aea4692938c2ad47b32a8c7b250dd9942c Mon Sep 17 00:00:00 2001 From: venkat1701 Date: Mon, 5 Oct 2026 05:17:01 +0530 Subject: [PATCH 3/3] docs(bench): describe crash recovery and what SIGKILL does not test --- docs/benchmarks.md | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 237f6d4..87f3e80 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -62,6 +62,7 @@ taken from a common hyperedge (so co-membership counts are non-zero), and 20,000 | `churn.remove` | remove one member from each of 20,000 other hyperedges, 1,000 per transaction | `Writer.remove` | `graph.replace` with a new `HGPlainLink` (links are immutable) | | `read.incidence.churned` | `read.incidence` after the deletes | | | | `disk.churned` | bytes on disk after the deletes and a clean shutdown | | | +| `recover` | open the store after the process ingesting into it was killed with SIGKILL | `HypergraphDatabase.open` | `HGEnvironment.get` | Read workloads run a warm-up pass over 10% of the probes before timing, except `read.incidence.cold`. @@ -92,6 +93,18 @@ the loaded store after a clean shutdown. The checksum of `read.incidence.churned both engines applied the same deletes. `disk.churned` shows whether the space came back: through compaction for HStore, and the log cleaner for JE. +**Crash recovery.** `reopen` closes the store cleanly first. `recover` doesn't: +1. The run starts a second JVM (`Comparison crash`) with the same flags and classpath. It ingests nodes into a + fresh directory and prints the count after every acknowledged commit. +2. Once half the nodes are acknowledged, the parent kills it with SIGKILL. +3. The parent times opening the directory, then counts the nodes that survived. + +The *Resources* table shows how many acknowledged nodes were lost, including when it's zero. SIGKILL is a process +crash, not a power failure: bytes the engine handed to the operating system survive it even without an fsync. So +this measures recovery time and the engine's own buffering, not whether `sync` really reaches the disk. Torn +writes and lost fsyncs are covered by HStore's crash matrix (`CrashRecoveryTest`, see [development.md](development.md)), not by this +benchmark. + **Bytes written.** For every write workload the report also shows the bytes each engine wrote per operation, in a *Resources* table, plus the total for the whole run. Both engines count the bytes they write themselves: