Skip to content

Commit 69380b7

Browse files
committed
feat(acp): implement lifecycle state machine
1 parent ecb8d9f commit 69380b7

1 file changed

Lines changed: 61 additions & 13 deletions

File tree

‎src/main/java/io/github/easy4j/kimi/acp/KimiAcpClient.java‎

Lines changed: 61 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
import java.util.concurrent.TimeUnit;
3535
import java.util.concurrent.atomic.AtomicBoolean;
3636
import java.util.concurrent.atomic.AtomicLong;
37+
import java.util.concurrent.atomic.AtomicReference;
3738
import java.util.function.Consumer;
3839

3940
import org.slf4j.Logger;
@@ -80,6 +81,8 @@ public class KimiAcpClient implements AutoCloseable {
8081
private final AtomicLong rpcIds = new AtomicLong();
8182
private final AtomicBoolean closed = new AtomicBoolean(false);
8283
private final AtomicBoolean connected = new AtomicBoolean(false);
84+
private final AtomicReference<KimiAcpState> state =
85+
new AtomicReference<KimiAcpState>(KimiAcpState.NEW);
8386
private final ScheduledExecutorService timer = Executors.newSingleThreadScheduledExecutor(r -> {
8487
Thread thread = new Thread(r, "kimi-acp-timer");
8588
thread.setDaemon(true);
@@ -113,10 +116,11 @@ public String connect() {
113116
if (closed.get()) {
114117
throw new IllegalStateException("kimi acp client is closed");
115118
}
116-
// CAS guard: a second connect would orphan the first child process.
117-
if (!connected.compareAndSet(false, true)) {
118-
throw new IllegalStateException("kimi acp client is already connected");
119+
if (!state.compareAndSet(KimiAcpState.NEW, KimiAcpState.CONNECTING)) {
120+
throw new IllegalStateException("kimi acp client is not connectable from state " + state.get());
119121
}
122+
// Keep the legacy connected flag aligned while callers migrate to getState().
123+
connected.set(true);
120124
List<String> command = new ArrayList<String>();
121125
command.add(config.getLocalExecutable());
122126
if (config.getAcpSubcommand() != null) {
@@ -133,8 +137,10 @@ public String connect() {
133137
process = builder.start();
134138
} catch (IOException e) {
135139
connected.set(false);
140+
state.set(KimiAcpState.NEW);
136141
throw new KimiException("Failed to spawn kimi acp: " + config.getLocalExecutable(), e);
137142
}
143+
state.set(KimiAcpState.INITIALIZING);
138144
stdin = new PrintWriter(new OutputStreamWriter(process.getOutputStream(), StandardCharsets.UTF_8), true);
139145
Thread reader = new Thread(this::readLoop, "kimi-acp-reader");
140146
reader.setDaemon(true);
@@ -152,6 +158,7 @@ public String connect() {
152158
if (result.hasNonNull("agentInfo")) {
153159
agentVersion = result.path("agentInfo").path("version").asText(null);
154160
}
161+
state.set(KimiAcpState.READY);
155162
return agentVersion;
156163
} catch (RuntimeException e) {
157164
// Handshake failure leaves the child alive — destroy it here so a
@@ -163,6 +170,7 @@ public String connect() {
163170
process = null;
164171
stdin = null;
165172
connected.set(false);
173+
state.set(KimiAcpState.NEW);
166174
throw e;
167175
}
168176
}
@@ -431,6 +439,15 @@ public boolean isClosed() {
431439
return closed.get();
432440
}
433441

442+
/**
443+
* Returns the current observable ACP lifecycle state.
444+
*
445+
* @return the current state; never {@code null}.
446+
*/
447+
public KimiAcpState getState() {
448+
return state.get();
449+
}
450+
434451
/**
435452
* Terminates the {@code kimi acp} child process and releases the timer.
436453
* Idempotent.
@@ -440,19 +457,23 @@ public void close() {
440457
if (!closed.compareAndSet(false, true)) {
441458
return;
442459
}
460+
state.set(KimiAcpState.CLOSING);
461+
connected.set(false);
443462
timer.shutdownNow();
444463
Process current = process;
445464
if (current != null) {
446465
current.destroy();
447466
}
448467
failAllPending(new KimiException("kimi acp client closed"));
468+
state.set(KimiAcpState.CLOSED);
449469
}
450470

451471
// ============================================================
452472
// transport internals
453473
// ============================================================
454474

455475
private CompletableFuture<JsonNode> request(String method, Map<String, Object> params) {
476+
requireRequestState(method);
456477
long id = rpcIds.incrementAndGet();
457478
Map<String, Object> payload = new LinkedHashMap<String, Object>();
458479
payload.put("jsonrpc", "2.0");
@@ -473,6 +494,7 @@ private CompletableFuture<JsonNode> request(String method, Map<String, Object> p
473494
}
474495

475496
private void notify(String method, Map<String, Object> params) {
497+
requireReady();
476498
Map<String, Object> payload = new LinkedHashMap<String, Object>();
477499
payload.put("jsonrpc", "2.0");
478500
payload.put("method", method);
@@ -517,16 +539,17 @@ private void readLoop() {
517539
KimiException error = new KimiException(
518540
"kimi acp frame exceeded maxFrameChars=" + config.getMaxFrameChars());
519541
log.warn("kimi acp frame over cap, tearing transport down");
520-
failAllPending(error);
521-
process.destroy();
542+
failTransport(error);
522543
return;
523544
}
524545
handleFrame(line);
525546
}
526-
failAllPending(new KimiException("kimi acp stdout closed (child exited)"));
547+
if (!closed.get()) {
548+
failTransport(new KimiException("kimi acp stdout closed (child exited)"));
549+
}
527550
} catch (IOException e) {
528551
if (!closed.get()) {
529-
failAllPending(new KimiException("kimi acp stdout read failed", e));
552+
failTransport(new KimiException("kimi acp stdout read failed", e));
530553
}
531554
}
532555
}
@@ -538,12 +561,7 @@ private void handleFrame(String frame) {
538561
} catch (Exception ex) {
539562
KimiException error = new KimiException("kimi acp protocol malformed JSON frame", ex);
540563
log.warn("kimi acp malformed JSON frame, tearing transport down");
541-
failAllPending(error);
542-
connected.set(false);
543-
Process current = process;
544-
if (current != null) {
545-
current.destroy();
546-
}
564+
failTransport(error);
547565
return;
548566
}
549567
if (node.hasNonNull("id")) {
@@ -616,6 +634,36 @@ private JsonNode await(CompletableFuture<JsonNode> future, long timeoutMillis, S
616634
}
617635
}
618636

637+
private void requireRequestState(String method) {
638+
KimiAcpState current = state.get();
639+
if ("initialize".equals(method) && current == KimiAcpState.INITIALIZING) {
640+
return;
641+
}
642+
if (current != KimiAcpState.READY) {
643+
throw new KimiException("kimi acp client is not ready: state=" + current);
644+
}
645+
}
646+
647+
private void requireReady() {
648+
KimiAcpState current = state.get();
649+
if (current != KimiAcpState.READY) {
650+
throw new KimiException("kimi acp client is not ready: state=" + current);
651+
}
652+
}
653+
654+
private void failTransport(KimiException error) {
655+
KimiAcpState currentState = state.get();
656+
connected.set(false);
657+
if (currentState == KimiAcpState.READY) {
658+
state.compareAndSet(KimiAcpState.READY, KimiAcpState.FAILED);
659+
}
660+
failAllPending(error);
661+
Process current = process;
662+
if (current != null) {
663+
current.destroy();
664+
}
665+
}
666+
619667
private void failAllPending(KimiException error) {
620668
for (Map.Entry<Long, CompletableFuture<JsonNode>> entry : pendingRpcs.entrySet()) {
621669
CompletableFuture<JsonNode> future = pendingRpcs.remove(entry.getKey());

0 commit comments

Comments
 (0)