Skip to content

Commit a845ed6

Browse files
committed
test: harden streaming and MCP recovery contracts
1 parent 08d5ceb commit a845ed6

5 files changed

Lines changed: 77 additions & 4 deletions

File tree

‎src/main/java/io/github/easy4j/comfy/cli/ComfyCli.java‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -394,9 +394,16 @@ public static class RunOptions {
394394
private boolean json;
395395
private boolean jsonStream;
396396

397+
public RunOptions() { }
398+
397399
public RunOptions(String workflowPath) {
398400
this.workflowPath = requireNonBlank("workflowPath", workflowPath);
399401
}
402+
403+
public RunOptions workflow(String value) {
404+
this.workflowPath = value == null ? null : requireNonBlank("workflowPath", value);
405+
return this;
406+
}
400407
public RunOptions wait(boolean value) { this.wait = Boolean.valueOf(value); return this; }
401408
public RunOptions prompt(String value) { this.prompt = value; return this; }
402409
public RunOptions set(String value) { this.setOverrides.add(requireNonBlank("set", value)); return this; }
@@ -414,6 +421,12 @@ public RunOptions(String workflowPath) {
414421
public RunOptions jsonStream(boolean value) { this.jsonStream = value; return this; }
415422

416423
List<String> toArgs(String defaultWhere) {
424+
if (workflowPath != null && prompt != null) {
425+
throw new IllegalStateException("run workflow and prompt are mutually exclusive");
426+
}
427+
if (workflowPath == null && (prompt == null || prompt.trim().isEmpty())) {
428+
throw new IllegalStateException("run requires workflow or prompt");
429+
}
417430
List<String> args = new ArrayList<String>();
418431
if (json) args.add("--json");
419432
if (jsonStream) args.add("--json-stream");

‎src/main/java/io/github/easy4j/comfy/cli/ComfyCliStreamExecutor.java‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -128,7 +128,7 @@ public ComfyCliStreamSession execute(final ComfyCliStreamListener listener, Stri
128128

129129
private void readStdout(Process process, ComfyCliStreamListener listener, CappedBytes retained,
130130
AtomicReference<Throwable> listenerFailure) {
131-
readLines(process.getInputStream(), new LineConsumer() {
131+
readLines(process.getInputStream(), effectiveLineLimit(config.getMaxStdoutBytes()), new LineConsumer() {
132132
@Override public void accept(String line) {
133133
retained.append((line + "\n").getBytes(StandardCharsets.UTF_8));
134134
if (line.trim().isEmpty()) return;
@@ -150,7 +150,7 @@ private void readStdout(Process process, ComfyCliStreamListener listener, Capped
150150

151151
private void readStderr(Process process, ComfyCliStreamListener listener, CappedBytes retained,
152152
AtomicReference<Throwable> listenerFailure) {
153-
readLines(process.getErrorStream(), new LineConsumer() {
153+
readLines(process.getErrorStream(), effectiveLineLimit(config.getMaxStderrBytes()), new LineConsumer() {
154154
@Override public void accept(String line) {
155155
retained.append((line + "\n").getBytes(StandardCharsets.UTF_8));
156156
try {
@@ -163,7 +163,7 @@ private void readStderr(Process process, ComfyCliStreamListener listener, Capped
163163
});
164164
}
165165

166-
private static void readLines(InputStream input, LineConsumer consumer) {
166+
private static void readLines(InputStream input, int maxLineBytes, LineConsumer consumer) {
167167
ByteArrayOutputStream line = new ByteArrayOutputStream();
168168
try {
169169
int b;
@@ -172,7 +172,7 @@ private static void readLines(InputStream input, LineConsumer consumer) {
172172
consumer.accept(new String(line.toByteArray(), StandardCharsets.UTF_8));
173173
line.reset();
174174
} else if (b != '\r') {
175-
line.write(b);
175+
if (line.size() < maxLineBytes) line.write(b);
176176
}
177177
}
178178
if (line.size() > 0) consumer.accept(new String(line.toByteArray(), StandardCharsets.UTF_8));
@@ -182,6 +182,10 @@ private static void readLines(InputStream input, LineConsumer consumer) {
182182
}
183183
}
184184

185+
private static int effectiveLineLimit(int configured) {
186+
return configured > 0 ? configured : 16 * 1024 * 1024;
187+
}
188+
185189
private void terminate(Process process) {
186190
if (!process.isAlive()) return;
187191
process.destroy();

‎src/test/java/io/github/easy4j/comfy/cli/ComfyCliParityTest.java‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,15 @@ void shouldValidateRoutingAndDynamicParameterNames() {
106106
() -> cli().skillsStatus("workspace"));
107107
}
108108

109+
@Test
110+
void runOptionsShouldSupportPromptOnlyAndRejectPromptWorkflowConflict() {
111+
String promptOnly = cli().run(new ComfyCli.RunOptions().prompt("a cat").noWatch(true))
112+
.getStdout();
113+
assertTrue(promptOnly.contains("run --prompt a cat --no-watch"));
114+
assertThrows(IllegalStateException.class,
115+
() -> cli().run(new ComfyCli.RunOptions("wf.json").prompt("cat")));
116+
}
117+
109118
@Test
110119
void copiedGenerateOptionsMustBeIndependent() {
111120
ComfyCli.GenerateOptions original = new ComfyCli.GenerateOptions().prompt("cat").param("quality", "high");

‎src/test/java/io/github/easy4j/comfy/mcp/ComfyMcpHardeningTest.java‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,44 @@ void notificationsMustBeSurfaced() {
9292
}
9393
}
9494

95+
@Test
96+
void elicitationShouldBeHandledOnlyWhenConfigured() {
97+
try (ComfyMcpClient client = new ComfyMcpClient(config())) {
98+
client.setElicitationHandler(new ComfyMcpElicitationHandler() {
99+
@Override public Object handle(String method, com.fasterxml.jackson.databind.JsonNode params) {
100+
Map<String, Object> result = new LinkedHashMap<String, Object>();
101+
result.put("accepted", Boolean.TRUE);
102+
return result;
103+
}
104+
});
105+
client.connect();
106+
ComfyMcpCallResult result = client.callTool("elicitation_test", null);
107+
assertFalse(result.isError());
108+
assertTrue(result.getText().contains("accepted=true"));
109+
}
110+
}
111+
112+
@Test
113+
void unexpectedChildExitShouldResetStateAndAllowReconnect() {
114+
ComfyMcpClient client = new ComfyMcpClient(config());
115+
try {
116+
client.connect();
117+
assertThrows(ComfyException.class, () -> client.callTool("die", null));
118+
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(2);
119+
while (client.isConnected() && System.nanoTime() < deadline) {
120+
try { Thread.sleep(10L); } catch (InterruptedException e) {
121+
Thread.currentThread().interrupt();
122+
break;
123+
}
124+
}
125+
assertFalse(client.isConnected());
126+
assertEquals("0.0.0-test", client.connect());
127+
assertTrue(client.isConnected());
128+
} finally {
129+
client.close();
130+
}
131+
}
132+
95133
@Test
96134
void nonTextMcpContentMustRemainStructured() {
97135
try (ComfyMcpClient client = new ComfyMcpClient(config())) {

‎src/test/resources/fake-mcp-server.py‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,15 @@ def main():
7171
send({"jsonrpc": "2.0", "method": "notifications/progress",
7272
"params": {"progress": 0.5}})
7373
reply(req_id, {"content": [{"type": "text", "text": "ok"}], "isError": False})
74+
elif name == "elicitation_test":
75+
send({"jsonrpc": "2.0", "id": 9001, "method": "elicitation/create",
76+
"params": {"message": "confirm"}})
77+
response = json.loads(sys.stdin.readline())
78+
accepted = bool((response.get("result") or {}).get("accepted"))
79+
reply(req_id, {"content": [{"type": "text", "text": "accepted=" + str(accepted).lower()}],
80+
"isError": False})
81+
elif name == "die":
82+
sys.exit(0)
7483
elif name == "nope":
7584
reply(req_id, {"content": [{"type": "text", "text": "unknown tool: " + name}],
7685
"isError": True})

0 commit comments

Comments
 (0)