Skip to content

Commit 409b7cc

Browse files
committed
feat(acp): expose lifecycle state and failed transport guard
1 parent 4ceb246 commit 409b7cc

2 files changed

Lines changed: 59 additions & 4 deletions

File tree

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

Lines changed: 37 additions & 4 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);
@@ -117,6 +120,7 @@ public String connect() {
117120
if (!connected.compareAndSet(false, true)) {
118121
throw new IllegalStateException("kimi acp client is already connected");
119122
}
123+
state.set(KimiAcpState.CONNECTING);
120124
List<String> command = new ArrayList<String>();
121125
command.add(config.getLocalExecutable());
122126
if (config.getAcpSubcommand() != null) {
@@ -131,8 +135,10 @@ public String connect() {
131135
builder.redirectErrorStream(false);
132136
try {
133137
process = builder.start();
138+
state.set(KimiAcpState.INITIALIZING);
134139
} catch (IOException e) {
135140
connected.set(false);
141+
state.set(KimiAcpState.NEW);
136142
throw new KimiException("Failed to spawn kimi acp: " + config.getLocalExecutable(), e);
137143
}
138144
stdin = new PrintWriter(new OutputStreamWriter(process.getOutputStream(), StandardCharsets.UTF_8), 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
}
@@ -436,6 +444,15 @@ public boolean isClosed() {
436444
return closed.get();
437445
}
438446

447+
/**
448+
* Returns the current ACP lifecycle state.
449+
*
450+
* @return the lifecycle state; never {@code null}.
451+
*/
452+
public KimiAcpState getState() {
453+
return state.get();
454+
}
455+
439456
/**
440457
* Terminates the {@code kimi acp} child process and releases the timer.
441458
* Idempotent.
@@ -445,6 +462,7 @@ public void close() {
445462
if (!closed.compareAndSet(false, true)) {
446463
return;
447464
}
465+
state.set(KimiAcpState.CLOSING);
448466
connected.set(false);
449467
timer.shutdownNow();
450468
PrintWriter writer = stdin;
@@ -458,13 +476,18 @@ public void close() {
458476
current.destroy();
459477
}
460478
failAllPending(new KimiException("kimi acp client closed"));
479+
state.set(KimiAcpState.CLOSED);
461480
}
462481

463482
// ============================================================
464483
// transport internals
465484
// ============================================================
466485

467486
private CompletableFuture<JsonNode> request(String method, Map<String, Object> params) {
487+
KimiAcpState currentState = state.get();
488+
if (!"initialize".equals(method) && currentState != KimiAcpState.READY) {
489+
throw new KimiException("kimi acp client is not ready: state=" + currentState);
490+
}
468491
long id = rpcIds.incrementAndGet();
469492
Map<String, Object> payload = new LinkedHashMap<String, Object>();
470493
payload.put("jsonrpc", "2.0");
@@ -530,20 +553,22 @@ private void readLoop() {
530553
KimiException error = new KimiException(
531554
"kimi acp frame exceeded maxFrameChars=" + config.getMaxFrameChars());
532555
log.warn("kimi acp frame over cap, tearing transport down");
533-
failAllPending(error);
556+
failTransport(error);
534557
process.destroy();
535558
return;
536559
}
537560
handleFrame(line);
538561
}
539-
failAllPending(new KimiException("kimi acp stdout closed (child exited)"));
562+
if (!closed.get()) {
563+
failTransport(new KimiException("kimi acp stdout closed (child exited)"));
564+
}
540565
} catch (IOException e) {
541566
if (!closed.get()) {
542-
failAllPending(new KimiException("kimi acp stdout read failed", e));
567+
failTransport(new KimiException("kimi acp stdout read failed", e));
543568
}
544569
} catch (KimiException e) {
545570
if (!closed.get()) {
546-
failAllPending(e);
571+
failTransport(e);
547572
Process current = process;
548573
if (current != null) {
549574
current.destroy();
@@ -631,6 +656,14 @@ private JsonNode await(CompletableFuture<JsonNode> future, long timeoutMillis, S
631656
}
632657
}
633658

659+
private void failTransport(KimiException error) {
660+
if (!closed.get()) {
661+
state.set(KimiAcpState.FAILED);
662+
}
663+
connected.set(false);
664+
failAllPending(error);
665+
}
666+
634667
private void failAllPending(KimiException error) {
635668
for (Map.Entry<Long, CompletableFuture<JsonNode>> entry : pendingRpcs.entrySet()) {
636669
CompletableFuture<JsonNode> future = pendingRpcs.remove(entry.getKey());
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
/*
2+
* Copyright (c) 2018-present, easy-4-java (https://github.com/easy-4-java).
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
*/
7+
package io.github.easy4j.kimi.acp;
8+
9+
/**
10+
* Observable lifecycle states of a {@link KimiAcpClient}.
11+
*
12+
* @since 3.0.0
13+
*/
14+
public enum KimiAcpState {
15+
NEW,
16+
CONNECTING,
17+
INITIALIZING,
18+
READY,
19+
CLOSING,
20+
CLOSED,
21+
FAILED
22+
}

0 commit comments

Comments
 (0)