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");
+ }
+}