Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
}

Expand All @@ -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();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* 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<JsonObject>(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<JsonObject> 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() {}
}
}
2 changes: 2 additions & 0 deletions keyext.api/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -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<JavaCompile>("compileJava") {
Expand Down
137 changes: 119 additions & 18 deletions keyext.api/src/main/java/org/keyproject/key/api/KeyApiImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,10 @@ public final class KeyApiImpl implements KeyApi {

private Function<Void, Boolean> 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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -552,25 +575,103 @@ public CompletableFuture<ProofId> loadExample(String name) {

@Override
public CompletableFuture<ProofId> 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<String> antecTerms, List<String> 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<String> 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<SortDesc> sorts) {
return String.join(", ", sorts.stream().map(SortDesc::string).toList());
}

private static void appendArgSorts(StringBuilder sb, List<SortDesc> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -16,7 +16,12 @@
* @version 1 (29.10.23)
*/
public class KeyIdentifications {
private final Map<EnvironmentId, KeyEnvironmentContainer> 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<EnvironmentId, KeyEnvironmentContainer> mapEnv =
new ConcurrentHashMap<>(16);

public KeyEnvironmentContainer getContainer(EnvironmentId environmentId) {
return Objects.requireNonNull(mapEnv.get(environmentId),
Expand Down Expand Up @@ -155,7 +160,7 @@ public record KeyEnvironmentContainer(KeYEnvironment<?> env,
Map<ProofId, ProofContainer> mapProof) {

public KeyEnvironmentContainer(KeYEnvironment<?> env) {
this(env, new HashMap<>(1));
this(env, new ConcurrentHashMap<>(1));
}

void dispose() {
Expand All @@ -172,7 +177,8 @@ private record ProofContainer(Proof wProof,
// where the position table is null) yields equal values.
Map<NodeTextId, NodeText> 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() {
Expand Down
Loading
Loading