diff --git a/keyext.api.client/src/main/java/org/key_project/key/api/client/RPCLayer.java b/keyext.api.client/src/main/java/org/key_project/key/api/client/RPCLayer.java index 3759323aef..e8206c136b 100644 --- a/keyext.api.client/src/main/java/org/key_project/key/api/client/RPCLayer.java +++ b/keyext.api.client/src/main/java/org/key_project/key/api/client/RPCLayer.java @@ -140,10 +140,14 @@ public void run() { if (message == null) break; final var jsonObject = JsonParser.parseString(message).getAsJsonObject(); - incomingMessageQueue.add(jsonObject); + // put() blocks when the queue is full (back-pressure); add() + // would throw IllegalStateException and kill this thread. + incomingMessageQueue.put(jsonObject); } } catch (IOException e) { throw new RuntimeException(e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); } } @@ -164,8 +168,15 @@ protected void processIncomingMessage(String message) { buf = new char[len]; } - int count = incoming.read(buf, 0, len); - assert count == len; + // read() may return fewer chars than requested (especially on a + // socket), so loop until the whole message has arrived or EOF. + int count = 0; + while (count < len) { + int n = incoming.read(buf, count, len - count); + if (n == -1) // EOF mid-message + break; + count += n; + } consumeCRNL(); return new String(buf, 0, count).trim(); } diff --git a/keyext.api.client/src/test/java/edu/kit/keyext/client/RPCLayerTest.java b/keyext.api.client/src/test/java/edu/kit/keyext/client/RPCLayerTest.java index 76fb5f4898..38b70b8a94 100644 --- a/keyext.api.client/src/test/java/edu/kit/keyext/client/RPCLayerTest.java +++ b/keyext.api.client/src/test/java/edu/kit/keyext/client/RPCLayerTest.java @@ -4,10 +4,13 @@ package edu.kit.keyext.client; import java.io.IOException; +import java.io.Reader; import java.io.StringReader; import java.io.StringWriter; import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.TimeUnit; +import com.google.gson.JsonObject; import org.key_project.key.api.client.JsonRPC; import org.key_project.key.api.client.RPCLayer; import org.junit.jupiter.api.Assertions; @@ -43,4 +46,104 @@ void testIncoming() throws IOException { String second = listener.readMessage(); Assertions.assertEquals(response, second); } + + /** + * Regression test for partial reads: a {@link Reader} (e.g. a socket) may + * return fewer characters than requested per {@code read} call. The listener + * must keep reading until the whole framed message has arrived; otherwise the + * message is silently truncated. + */ + @Test + void testIncomingWithPartialReads() throws IOException { + var response = JsonRPC.createResponse("1", 2); + var notification = JsonRPC.createNotification("test", "some longer payload"); + var incoming = new OneCharAtATimeReader( + JsonRPC.addHeader(response) + JsonRPC.addHeader(notification)); + RPCLayer.JsonStreamListener listener = + new RPCLayer.JsonStreamListener(incoming, new ArrayBlockingQueue<>(1)); + Assertions.assertEquals(response, listener.readMessage()); + Assertions.assertEquals(notification, listener.readMessage()); + Assertions.assertNull(listener.readMessage()); // clean EOF + } + + /** + * Regression test for back-pressure: when the queue is full the listener must + * block ({@code put}) rather than throw ({@code add}). With the old + * {@code add()} the second message threw {@link IllegalStateException} and + * killed the reader thread, so only the first message was ever delivered. + *

+ * The check is made deterministic by not draining the queue until the + * listener has had to insert the second message into the full (capacity-1) + * queue: with {@code put} it parks (thread stays alive), with {@code add} it + * has already died. + */ + @Test + void testListenerBlocksWhenQueueFull() throws Exception { + var m1 = JsonRPC.addHeader(JsonRPC.createResponse("1", 1)); + var m2 = JsonRPC.addHeader(JsonRPC.createResponse("2", 2)); + var m3 = JsonRPC.addHeader(JsonRPC.createResponse("3", 3)); + var queue = new ArrayBlockingQueue(1); + var listener = new RPCLayer.JsonStreamListener(new StringReader(m1 + m2 + m3), queue); + + var thread = new Thread(listener, "test-listener"); + thread.setDaemon(true); + thread.start(); + + // Wait until the first message is queued (slot full) and the listener is + // attempting the second one, i.e. it has either parked (put) or died (add). + awaitUntil(() -> queue.size() == 1); + awaitUntil(() -> !thread.isAlive() + || thread.getState() == Thread.State.WAITING + || thread.getState() == Thread.State.TIMED_WAITING); + Assertions.assertTrue(thread.isAlive(), + "listener died instead of applying back-pressure on a full queue"); + + // Draining now lets the blocked put()s proceed; all three arrive in order. + Assertions.assertEquals(1, take(queue).get("result").getAsInt()); + Assertions.assertEquals(2, take(queue).get("result").getAsInt()); + Assertions.assertEquals(3, take(queue).get("result").getAsInt()); + + thread.join(5000); + Assertions.assertFalse(thread.isAlive(), "listener thread should finish at EOF"); + } + + private static JsonObject take(ArrayBlockingQueue queue) + throws InterruptedException { + var v = queue.poll(5, TimeUnit.SECONDS); + Assertions.assertNotNull(v, "expected a message but none arrived"); + return v; + } + + private static void awaitUntil(java.util.function.BooleanSupplier condition) + throws InterruptedException { + long deadline = System.currentTimeMillis() + 5000; + while (!condition.getAsBoolean() && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + } + + /** A {@link Reader} that hands out at most one character per read call. */ + private static final class OneCharAtATimeReader extends Reader { + private final String data; + private int pos = 0; + + OneCharAtATimeReader(String data) { + this.data = data; + } + + @Override + public int read(char[] cbuf, int off, int len) { + if (pos >= data.length()) { + return -1; + } + if (len <= 0) { + return 0; + } + cbuf[off] = data.charAt(pos++); + return 1; + } + + @Override + public void close() {} + } } diff --git a/keyext.api/build.gradle.kts b/keyext.api/build.gradle.kts index a6c60d3313..a013fda955 100644 --- a/keyext.api/build.gradle.kts +++ b/keyext.api/build.gradle.kts @@ -23,6 +23,8 @@ dependencies { implementation(libs.guava) annotationProcessor(libs.therapi.runtime.javadoc.scribe) api(libs.therapi.runtime.javadoc) + + testImplementation("org.mockito:mockito-core:5.18.0") } tasks.named("compileJava") { diff --git a/keyext.api/src/main/java/org/keyproject/key/api/KeyApiImpl.java b/keyext.api/src/main/java/org/keyproject/key/api/KeyApiImpl.java index 105ffdae8b..324b4830c5 100644 --- a/keyext.api/src/main/java/org/keyproject/key/api/KeyApiImpl.java +++ b/keyext.api/src/main/java/org/keyproject/key/api/KeyApiImpl.java @@ -67,7 +67,10 @@ public final class KeyApiImpl implements KeyApi { private Function exitHandler; - private ClientApi clientApi; + /// Never {@code null}: starts as a no-op sink so that task/loading callbacks + /// fired before a client has connected (or after it disconnects) cannot hit + /// a null reference. {@link #setClientApi} swaps in the real remote proxy. + private ClientApi clientApi = noopClientApi(); private final ProverTaskListener clientListener = new ProverTaskListener() { @Override public void taskStarted(org.key_project.prover.engine.TaskStartedInfo info) { @@ -516,6 +519,26 @@ public void setClientApi(ClientApi remoteProxy) { clientApi = remoteProxy; } + /// Detaches the current client, e.g. after it disconnects in TCP server mode. + /// Subsequent notifications go to the no-op sink instead of a dead remote + /// proxy, until a new client attaches via {@link #setClientApi}. + public void disconnectClient() { + clientApi = noopClientApi(); + } + + /// A {@link ClientApi} whose calls do nothing: notifications are dropped and + /// requests resolve to {@code null}. Used as the default before a client is + /// attached, see {@link #clientApi}. + private static ClientApi noopClientApi() { + return (ClientApi) java.lang.reflect.Proxy.newProxyInstance( + ClientApi.class.getClassLoader(), + new Class[] { ClientApi.class }, + (proxy, method, args) -> CompletableFuture.class + .isAssignableFrom(method.getReturnType()) + ? CompletableFuture.completedFuture(null) + : null); + } + private final DefaultUserInterfaceControl control = new MyDefaultUserInterfaceControl(); @Override @@ -552,25 +575,103 @@ public CompletableFuture loadExample(String name) { @Override public CompletableFuture loadProblem(ProblemDefinition problem) { - return CompletableFutures.computeAsync((c) -> { - Proof proof = null; - KeYEnvironment env = null; - /* - * var loader = control.load(JavaProfile.getDefaultProfile(), - * ex.getObligationFile(), null, null, null, null, true, null); - * InitConfig initConfig = loader.getInitConfig(); - * - * env = new KeYEnvironment<>(control, initConfig, loader.getProof(), - * loader.getProofScript(), loader.getResult()); - * var envId = new EnvironmentId(env.toString()); - * data.register(envId, env); - * proof = Objects.requireNonNull(env.getLoadedProof()); - * var proofId = new ProofId(envId, proof.name().toString()); - * return data.register(proofId, proof); - */ + // Render the problem definition into a KeY input file and load it through + // the regular .key loading path. + return loadKey(buildKeyInput(problem)); + } + + /// Builds a {@code .key} input document from a {@link ProblemDefinition}: the + /// declared sorts/functions/predicates followed by a {@code \problem} holding + /// the sequent {@code antecTerms ==> succTerms} encoded as a single formula. + static String buildKeyInput(ProblemDefinition problem) { + var sb = new StringBuilder(); + + var sorts = problem.sorts(); + if (sorts != null && !sorts.isEmpty()) { + sb.append("\\sorts {\n"); + for (var sort : sorts) { + sb.append(" "); + if (sort.anAbstract()) { + sb.append("\\abstract "); + } + sb.append(sort.string()); + var ext = sort.extendsSorts(); + if (ext != null && !ext.isEmpty()) { + sb.append(" \\extends ").append(joinSortNames(ext)); + } + sb.append(";\n"); + } + sb.append("}\n\n"); + } + + var functions = problem.functions(); + if (functions != null && !functions.isEmpty()) { + sb.append("\\functions {\n"); + for (var fn : functions) { + sb.append(" ").append(returnSortName(fn)).append(" ").append(fn.name()); + appendArgSorts(sb, fn.argSorts()); + sb.append(";\n"); + } + sb.append("}\n\n"); + } + + var predicates = problem.predicates(); + if (predicates != null && !predicates.isEmpty()) { + sb.append("\\predicates {\n"); + for (var pred : predicates) { + sb.append(" ").append(pred.name()); + appendArgSorts(sb, pred.argSorts()); + sb.append(";\n"); + } + sb.append("}\n\n"); + } + + sb.append("\\problem {\n ") + .append(buildSequentFormula(problem.antecTerms(), problem.succTerms())) + .append("\n}\n"); + return sb.toString(); + } + + /// Encodes the sequent {@code antec ==> succ} as a single formula + /// {@code (a1 & ... & an) -> (s1 | ... | sm)}. An empty succedent becomes + /// {@code false}; an empty antecedent drops the implication. + private static String buildSequentFormula(List antecTerms, List succTerms) { + String ante = joinFormulas(antecTerms, " & "); + String succ = joinFormulas(succTerms, " | "); + if (ante == null && succ == null) { + throw new IllegalArgumentException( + "ProblemDefinition must contain at least one antecedent or succedent term"); + } + if (ante == null) { + return succ; + } + if (succ == null) { + return ante + " -> false"; + } + return ante + " -> " + succ; + } + + /// Joins terms with {@code op}, parenthesising each, or {@code null} if empty. + private static String joinFormulas(List terms, String op) { + if (terms == null || terms.isEmpty()) { return null; - }); + } + return String.join(op, terms.stream().map(t -> "(" + t + ")").toList()); + } + + private static String joinSortNames(List sorts) { + return String.join(", ", sorts.stream().map(SortDesc::string).toList()); + } + + private static void appendArgSorts(StringBuilder sb, List argSorts) { + if (argSorts == null || argSorts.isEmpty()) { + return; + } + sb.append("(").append(joinSortNames(argSorts)).append(")"); + } + private static String returnSortName(FunctionDesc fn) { + return fn.retSort() != null ? fn.retSort().string() : fn.sort(); } @Override diff --git a/keyext.api/src/main/java/org/keyproject/key/api/StartServer.java b/keyext.api/src/main/java/org/keyproject/key/api/StartServer.java index 4d30491f60..93aaeafc25 100644 --- a/keyext.api/src/main/java/org/keyproject/key/api/StartServer.java +++ b/keyext.api/src/main/java/org/keyproject/key/api/StartServer.java @@ -227,6 +227,9 @@ private void runTcpServer(KeyApiImpl keyApi) } catch (IOException e) { LOGGER.warn("Connection error; awaiting next connection", e); } + // Detach the gone client so notifications don't hit a dead proxy + // before the next client connects. + keyApi.disconnectClient(); LOGGER.info("Client disconnected; awaiting next connection"); } } diff --git a/keyext.api/src/main/java/org/keyproject/key/api/data/KeyIdentifications.java b/keyext.api/src/main/java/org/keyproject/key/api/data/KeyIdentifications.java index 9174d0cc84..7d5993ac95 100644 --- a/keyext.api/src/main/java/org/keyproject/key/api/data/KeyIdentifications.java +++ b/keyext.api/src/main/java/org/keyproject/key/api/data/KeyIdentifications.java @@ -3,9 +3,9 @@ * SPDX-License-Identifier: GPL-2.0-only */ package org.keyproject.key.api.data; -import java.util.HashMap; import java.util.Map; import java.util.Objects; +import java.util.concurrent.ConcurrentHashMap; import de.uka.ilkd.key.control.KeYEnvironment; import de.uka.ilkd.key.proof.Node; @@ -16,7 +16,12 @@ * @version 1 (29.10.23) */ public class KeyIdentifications { - private final Map mapEnv = new HashMap<>(16); + // Concurrent: the JSON-RPC launcher dispatches requests on multiple threads, + // so register/find/dispose can run in parallel. A plain HashMap would risk + // lost updates and ConcurrentModificationException (e.g. dispose(env) iterates + // mapProof while dispose(proof) removes from it). + private final Map mapEnv = + new ConcurrentHashMap<>(16); public KeyEnvironmentContainer getContainer(EnvironmentId environmentId) { return Objects.requireNonNull(mapEnv.get(environmentId), @@ -155,7 +160,7 @@ public record KeyEnvironmentContainer(KeYEnvironment env, Map mapProof) { public KeyEnvironmentContainer(KeYEnvironment env) { - this(env, new HashMap<>(1)); + this(env, new ConcurrentHashMap<>(1)); } void dispose() { @@ -172,7 +177,8 @@ private record ProofContainer(Proof wProof, // where the position table is null) yields equal values. Map mapGoalText) { public ProofContainer(Proof proof) { - this(proof, new HashMap<>(16), new HashMap<>(16), new HashMap<>(16)); + this(proof, new ConcurrentHashMap<>(16), new ConcurrentHashMap<>(16), + new ConcurrentHashMap<>(16)); } void dispose() { diff --git a/keyext.api/src/test/java/org/keyproject/key/api/ClientDisconnectTest.java b/keyext.api/src/test/java/org/keyproject/key/api/ClientDisconnectTest.java new file mode 100644 index 0000000000..f8bcb71b49 --- /dev/null +++ b/keyext.api/src/test/java/org/keyproject/key/api/ClientDisconnectTest.java @@ -0,0 +1,34 @@ +/* This file is part of KeY - https://key-project.org + * KeY is licensed under the GNU General Public License Version 2 + * SPDX-License-Identifier: GPL-2.0-only */ +package org.keyproject.key.api; + +import java.lang.reflect.Field; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.keyproject.key.api.remoteclient.ClientApi; + +import static org.mockito.Mockito.mock; + +/** + * Test for reconnection handling: after a client disconnects, the API must + * detach it (fall back to the no-op sink) rather than keep a dead remote proxy. + */ +class ClientDisconnectTest { + @Test + void disconnectClientDetachesTheClient() throws Exception { + var api = new KeyApiImpl(); + var client = mock(ClientApi.class); + api.setClientApi(client); + + Field field = KeyApiImpl.class.getDeclaredField("clientApi"); + field.setAccessible(true); + Assertions.assertSame(client, field.get(api), "client should be attached"); + + api.disconnectClient(); + + Assertions.assertNotSame(client, field.get(api), "client should be detached"); + Assertions.assertNotNull(field.get(api), "should fall back to a no-op, not null"); + } +} diff --git a/keyext.api/src/test/java/org/keyproject/key/api/KeyApiImplTest.java b/keyext.api/src/test/java/org/keyproject/key/api/KeyApiImplTest.java new file mode 100644 index 0000000000..a3c498ba2c --- /dev/null +++ b/keyext.api/src/test/java/org/keyproject/key/api/KeyApiImplTest.java @@ -0,0 +1,102 @@ +/* This file is part of KeY - https://key-project.org + * KeY is licensed under the GNU General Public License Version 2 + * SPDX-License-Identifier: GPL-2.0-only */ +package org.keyproject.key.api; + +import java.util.List; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.keyproject.key.api.data.FunctionDesc; +import org.keyproject.key.api.data.PredicateDesc; +import org.keyproject.key.api.data.PrintOptions; +import org.keyproject.key.api.data.ProblemDefinition; +import org.keyproject.key.api.data.SortDesc; + +/** + * Tests for {@link KeyApiImpl#loadProblem} and its {@code .key} rendering. + */ +class KeyApiImplTest { + private static SortDesc sort(String name) { + return new SortDesc(name, "", List.of(), false, name); + } + + /** Pure rendering test: no KeY loading, just the generated input document. */ + @Test + void buildKeyInputRendersDeclarationsAndSequent() { + var s = sort("s"); + var problem = new ProblemDefinition( + List.of(s), + List.of(new FunctionDesc("c", "s", s, List.of(), true, false, false)), + List.of(new PredicateDesc("p", List.of(s))), + List.of("p(c)"), + List.of("p(c)")); + + var key = KeyApiImpl.buildKeyInput(problem); + Assertions.assertTrue(key.contains("\\sorts {"), key); + Assertions.assertTrue(key.contains("s;"), key); + Assertions.assertTrue(key.contains("\\functions {"), key); + Assertions.assertTrue(key.contains("s c;"), key); + Assertions.assertTrue(key.contains("\\predicates {"), key); + Assertions.assertTrue(key.contains("p(s);"), key); + Assertions.assertTrue(key.contains("\\problem {"), key); + Assertions.assertTrue(key.contains("(p(c)) -> (p(c))"), key); + } + + @Test + void buildKeyInputEncodesEmptyAntecedentOrSuccedent() { + // succedent only -> the disjunction of goals, no implication + var succOnly = new ProblemDefinition(null, null, null, null, List.of("a", "b")); + Assertions.assertTrue(KeyApiImpl.buildKeyInput(succOnly).contains("(a) | (b)")); + // antecedent only -> assumptions imply false + var anteOnly = new ProblemDefinition(null, null, null, List.of("a"), null); + Assertions.assertTrue(KeyApiImpl.buildKeyInput(anteOnly).contains("(a) -> false")); + // neither -> nothing to prove + var empty = new ProblemDefinition(null, null, null, null, null); + Assertions.assertThrows(IllegalArgumentException.class, + () -> KeyApiImpl.buildKeyInput(empty)); + } + + /** + * End-to-end: a programmatic problem with a custom sort, constant, and two + * predicates is loaded, and the resulting root sequent is the expected + * {@code p(c) ==> q(c)}. + */ + @Test + void loadProblemLoadsTheDefinedSequent() throws Exception { + var s = sort("s"); + var problem = new ProblemDefinition( + List.of(s), + List.of(new FunctionDesc("c", "s", s, List.of(), true, false, false)), + List.of(new PredicateDesc("p", List.of(s)), new PredicateDesc("q", List.of(s))), + List.of("p(c)"), + List.of("q(c)")); + + var api = new KeyApiImpl(); + var proofId = api.loadProblem(problem).get(); + Assertions.assertNotNull(proofId, "expected a proof to be created"); + + var root = api.root(proofId).get(); + var printed = api.print(root.nodeid(), new PrintOptions(false, 80, 2, false, false)).get(); + var sequent = printed.sequent(); + Assertions.assertTrue(sequent.contains("p(c)"), sequent); + Assertions.assertTrue(sequent.contains("q(c)"), sequent); + Assertions.assertTrue(sequent.contains("==>"), sequent); + } + + /** Abstract sorts and {@code \extends} are rendered and parsed correctly. */ + @Test + void loadProblemSupportsAbstractSortsAndSubsorts() throws Exception { + var a = new SortDesc("a", "", List.of(), true, "a"); // abstract + var b = new SortDesc("b", "", List.of(a), false, "b"); // b extends a + var problem = new ProblemDefinition( + List.of(a, b), + List.of(new FunctionDesc("c", "b", b, List.of(), true, false, false)), + List.of(new PredicateDesc("p", List.of(a))), + List.of("p(c)"), + List.of("p(c)")); + + var api = new KeyApiImpl(); + Assertions.assertNotNull(api.loadProblem(problem).get()); + } +} diff --git a/keyext.api/src/test/java/org/keyproject/key/api/data/KeyIdentificationsTest.java b/keyext.api/src/test/java/org/keyproject/key/api/data/KeyIdentificationsTest.java new file mode 100644 index 0000000000..bbca21338e --- /dev/null +++ b/keyext.api/src/test/java/org/keyproject/key/api/data/KeyIdentificationsTest.java @@ -0,0 +1,102 @@ +/* This file is part of KeY - https://key-project.org + * KeY is licensed under the GNU General Public License Version 2 + * SPDX-License-Identifier: GPL-2.0-only */ +package org.keyproject.key.api.data; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import de.uka.ilkd.key.control.KeYEnvironment; +import de.uka.ilkd.key.proof.Proof; + +import org.key_project.logic.Name; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.keyproject.key.api.data.KeyIdentifications.EnvironmentId; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Concurrency test for {@link KeyIdentifications}. The JSON-RPC launcher + * dispatches requests on several threads, so register/query must be thread-safe. + * With a plain {@link java.util.HashMap} this test throws + * {@link java.util.ConcurrentModificationException} or loses updates; with + * concurrent maps it is stable. + */ +class KeyIdentificationsTest { + @Test + void concurrentRegisterAndAllProofIdsIsThreadSafe() throws Exception { + var data = new KeyIdentifications(); + var envId = new EnvironmentId("env"); + data.register(envId, mock(KeYEnvironment.class)); + + final int writers = 8; + final int perWriter = 150; + + // Pre-build the proof mocks: mock creation is not what we are testing, + // and keeping it out of the concurrent section isolates the map races. + List> batches = new ArrayList<>(); + for (int w = 0; w < writers; w++) { + List batch = new ArrayList<>(perWriter); + for (int j = 0; j < perWriter; j++) { + var proof = mock(Proof.class); + when(proof.name()).thenReturn(new Name("p_" + w + "_" + j)); + batch.add(proof); + } + batches.add(batch); + } + + var pool = Executors.newFixedThreadPool(writers + 2); + var start = new CountDownLatch(1); + var errors = new CopyOnWriteArrayList(); + var tasks = new ArrayList>(); + + // Writers register proofs into the same environment concurrently. + for (List batch : batches) { + tasks.add(() -> { + start.await(); + for (Proof proof : batch) { + data.register(envId, proof); + } + return null; + }); + } + // Readers iterate the proof map concurrently with the writers. + for (int r = 0; r < 2; r++) { + tasks.add(() -> { + start.await(); + for (int k = 0; k < 1000; k++) { + data.allProofIds(); + } + return null; + }); + } + + List> futures = new ArrayList<>(); + for (Callable task : tasks) { + futures.add(pool.submit(task)); + } + start.countDown(); + for (Future f : futures) { + try { + f.get(30, TimeUnit.SECONDS); + } catch (ExecutionException e) { + errors.add(e.getCause()); + } + } + pool.shutdownNow(); + + Assertions.assertTrue(errors.isEmpty(), () -> "concurrent access failed: " + errors); + Assertions.assertEquals(writers * perWriter, data.allProofIds().size(), + "every distinct proof id should be registered exactly once"); + } +}