Skip to content

Commit 9a86100

Browse files
committed
feat(runtime): stream CLI output before process exit
1 parent 7231e70 commit 9a86100

1 file changed

Lines changed: 128 additions & 0 deletions

File tree

‎src/main/java/io/github/easy4j/opencode/cli/OpenCodeCliExecutor.java‎

Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,14 +9,22 @@
99
import org.slf4j.LoggerFactory;
1010

1111
import java.io.ByteArrayOutputStream;
12+
import java.io.BufferedReader;
1213
import java.io.File;
1314
import java.io.IOException;
15+
import java.io.InputStream;
16+
import java.io.InputStreamReader;
1417
import java.io.OutputStream;
1518
import java.nio.charset.StandardCharsets;
19+
import java.util.ArrayList;
20+
import java.util.Collections;
1621
import java.util.HashMap;
22+
import java.util.List;
1723
import java.util.Map;
1824
import java.util.Objects;
25+
import java.util.concurrent.CompletableFuture;
1926
import java.util.concurrent.Semaphore;
27+
import java.util.function.Consumer;
2028

2129
/**
2230
* Executor for the local {@code opencode} CLI subprocess.
@@ -90,6 +98,126 @@ public OpenCodeCliResult execute(OpenCodeCliExecutionContext context, String...
9098
}
9199
}
92100

101+
/**
102+
* Start a child process and deliver complete UTF-8 stdout/stderr lines while it is running.
103+
* The returned completion future resolves after process exit and stream pumps finish.
104+
*/
105+
public OpenCodeCliStreamHandle stream(OpenCodeCliExecutionContext context,
106+
Consumer<String> stdoutConsumer,
107+
Consumer<String> stderrConsumer,
108+
String... args) {
109+
boolean acquired = false;
110+
try {
111+
if (concurrencyLimiter != null) {
112+
concurrencyLimiter.acquire();
113+
acquired = true;
114+
}
115+
116+
ProcessBuilder builder = createProcessBuilder(context, args);
117+
Process process = builder.start();
118+
CompletableFuture<OpenCodeCliResult> completion = new CompletableFuture<>();
119+
OpenCodeCliStreamHandle handle = new OpenCodeCliStreamHandle(process, completion);
120+
121+
BoundedOutputStream stdout = new BoundedOutputStream(config.getMaxStdoutBytes());
122+
BoundedOutputStream stderr = new BoundedOutputStream(config.getMaxStderrBytes());
123+
124+
Thread stdoutThread = pumpThread(
125+
"opencode-cli-stdout", process.getInputStream(), stdout, stdoutConsumer);
126+
Thread stderrThread = pumpThread(
127+
"opencode-cli-stderr", process.getErrorStream(), stderr, stderrConsumer);
128+
129+
final boolean releasePermit = acquired;
130+
Thread waiter = new Thread(() -> {
131+
try {
132+
int exitCode = process.waitFor();
133+
joinPump(stdoutThread);
134+
joinPump(stderrThread);
135+
completion.complete(buildResult(exitCode, stdout, stderr, null));
136+
} catch (InterruptedException interrupted) {
137+
Thread.currentThread().interrupt();
138+
handle.cancel();
139+
completion.complete(buildResult(-1, stdout, stderr, "CLI streaming wait interrupted"));
140+
} finally {
141+
if (releasePermit && concurrencyLimiter != null) {
142+
concurrencyLimiter.release();
143+
}
144+
}
145+
}, "opencode-cli-waiter");
146+
waiter.setDaemon(true);
147+
waiter.start();
148+
149+
return handle;
150+
} catch (InterruptedException interrupted) {
151+
Thread.currentThread().interrupt();
152+
CompletableFuture<OpenCodeCliResult> failed = CompletableFuture.completedFuture(
153+
new OpenCodeCliResult(-1, "", "CLI streaming start interrupted"));
154+
return new OpenCodeCliStreamHandle(null, failed);
155+
} catch (IOException error) {
156+
if (acquired && concurrencyLimiter != null) {
157+
concurrencyLimiter.release();
158+
}
159+
CompletableFuture<OpenCodeCliResult> failed = CompletableFuture.completedFuture(
160+
new OpenCodeCliResult(-1, "", error.getMessage()));
161+
return new OpenCodeCliStreamHandle(null, failed);
162+
}
163+
}
164+
165+
public OpenCodeCliStreamHandle stream(Consumer<String> stdoutConsumer,
166+
Consumer<String> stderrConsumer,
167+
String... args) {
168+
return stream(null, stdoutConsumer, stderrConsumer, args);
169+
}
170+
171+
private ProcessBuilder createProcessBuilder(OpenCodeCliExecutionContext context, String... args) {
172+
List<String> command = new ArrayList<>();
173+
command.add(config.getExecutable());
174+
Collections.addAll(command, args);
175+
176+
ProcessBuilder builder = new ProcessBuilder(command);
177+
File workingDirectory = resolveWorkingDirectory(context);
178+
if (workingDirectory != null) {
179+
builder.directory(workingDirectory);
180+
}
181+
182+
if (context != null) {
183+
Map<String, String> environment = builder.environment();
184+
if (!context.isInheritParentEnvironment()) {
185+
environment.clear();
186+
}
187+
environment.putAll(context.getEnvironment());
188+
}
189+
return builder;
190+
}
191+
192+
private Thread pumpThread(String name, InputStream input,
193+
BoundedOutputStream capture, Consumer<String> consumer) {
194+
Thread thread = new Thread(() -> {
195+
try (BufferedReader reader = new BufferedReader(
196+
new InputStreamReader(input, StandardCharsets.UTF_8))) {
197+
String line;
198+
while ((line = reader.readLine()) != null) {
199+
byte[] bytes = line.getBytes(StandardCharsets.UTF_8);
200+
capture.write(bytes, 0, bytes.length);
201+
capture.write('\n');
202+
if (consumer != null) {
203+
consumer.accept(line);
204+
}
205+
}
206+
} catch (IOException error) {
207+
if (config.getDebug().allows(HttpLogLevel.BASIC)) {
208+
log.debug("OpenCode CLI stream pump ended: stream={}, error={}", name, error.getMessage());
209+
}
210+
}
211+
}, name);
212+
thread.setDaemon(true);
213+
thread.start();
214+
return thread;
215+
}
216+
217+
private void joinPump(Thread thread) throws InterruptedException {
218+
thread.join();
219+
}
220+
93221
private OpenCodeCliResult executeInternal(OpenCodeCliExecutionContext context, String... args) {
94222
CommandLine cmd = CommandLine.parse(config.getExecutable());
95223
for (String arg : args) {

0 commit comments

Comments
 (0)