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
141 changes: 114 additions & 27 deletions java/src/main/java/dev/sorted/mcphub/ControlHandler.java
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,14 @@
import org.slf4j.LoggerFactory;

import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

/**
Expand Down Expand Up @@ -40,6 +47,9 @@ public class ControlHandler implements JsonRpcServer.MethodHandler {
private BodyBudgetService bodyBudget;
private ProviderHealthTracker healthTracker;
private final AtomicInteger activeBridgeCount = new AtomicInteger(0);
private final Set<Long> bridgePids = ConcurrentHashMap.newKeySet();
private final Map<Long, Instant> bridgeLastPing = new ConcurrentHashMap<>();
private final ScheduledExecutorService janitorExecutor;

/** Backward-compatible constructor (Session 1 tests). */
public ControlHandler(StateMachine stateMachine, SessionManager sessionManager,
Expand All @@ -51,7 +61,13 @@ public ControlHandler(StateMachine stateMachine, SessionManager sessionManager,
this.policy = null;
this.bodyBudget = null;
this.healthTracker = null;
this.janitorExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
Thread t = new Thread(r, "mcphub-bridge-janitor");
t.setDaemon(true);
return t;
});
wireCallbacks();
janitorExecutor.scheduleWithFixedDelay(this::janitorTask, 60, 60, TimeUnit.SECONDS);
}

/** Full constructor for Session 2+. */
Expand All @@ -65,7 +81,18 @@ public ControlHandler(StateMachine stateMachine, SessionManager sessionManager,
this.policy = policy;
this.bodyBudget = bodyBudget;
this.healthTracker = null;
this.janitorExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
Thread t = new Thread(r, "mcphub-bridge-janitor");
t.setDaemon(true);
return t;
});
wireCallbacks();
janitorExecutor.scheduleWithFixedDelay(this::janitorTask, 60, 60, TimeUnit.SECONDS);
}

/** Shutdown the bridge janitor executor. Call on daemon shutdown or test teardown. */
public void shutdown() {
janitorExecutor.shutdownNow();
}

/** Optional wiring for REQ-4.7.2 runtime provider health visibility. */
Expand Down Expand Up @@ -103,13 +130,11 @@ private void wireCallbacks() {
if (healthTracker != null && to == StateMachine.State.COOLING_DOWN) {
healthTracker.clear();
}
// Track active bridge count for last-bridge-exit auto-close
if (from == StateMachine.State.CLOSED && to != StateMachine.State.CLOSED) {
activeBridgeCount.incrementAndGet();
log.info("Session started, active bridges: {}", activeBridgeCount.get());
}
// Defensive cleanup when session fully closes
if (to == StateMachine.State.CLOSED) {
activeBridgeCount.set(0);
bridgePids.clear();
bridgeLastPing.clear();
log.info("Session closed, active bridges reset to 0");
}
});
Expand All @@ -118,6 +143,10 @@ private void wireCallbacks() {
sessionManager.setTimeoutCallback(new SessionManager.TimeoutCallback() {
@Override
public void onIdleTimeout(String sessionId) {
if (activeBridgeCount.get() > 0) {
log.debug("Idle timeout suppressed: {} bridges attached", activeBridgeCount.get());
return;
}
log.info("Idle timeout for session {}", sessionId);
try {
stateMachine.transition(StateMachine.Trigger.IDLE_TIMEOUT, sessionId);
Expand Down Expand Up @@ -156,7 +185,9 @@ public JsonNode handle(String method, JsonNode params)
case "mcphub.control.arm" -> handleArm();
case "mcphub.control.open" -> handleOpen();
case "mcphub.control.close" -> handleClose();
case "mcphub.control.bridge_detach" -> handleBridgeDetach();
case "mcphub.control.bridge_attach"-> handleBridgeAttach(params);
case "mcphub.control.bridge_ping" -> handleBridgePing(params);
case "mcphub.control.bridge_detach"-> handleBridgeDetach(params);
case "mcphub.control.lock" -> handleLock(params);
case "mcphub.control.unlock" -> handleUnlock();
case "mcphub.control.health" -> handleHealth();
Expand Down Expand Up @@ -365,34 +396,63 @@ private JsonNode handleCapabilities() {
return r;
}

/** mcphub.control.bridge_detach — decrement bridge count; close session when last bridge leaves. */
private JsonNode handleBridgeDetach() {
int count = activeBridgeCount.updateAndGet(c -> c > 0 ? c - 1 : 0);
log.info("Bridge detached. Active bridges: {}", count);
if (count == 0) {
String sessionId = sessionManager.getCurrentSessionId();
StateMachine.State current = stateMachine.getState();
try {
if (current == StateMachine.State.ARMED) {
stateMachine.transition(StateMachine.Trigger.CLOSE, sessionId);
sessionManager.endSession();
} else if (current == StateMachine.State.OPEN) {
stateMachine.transition(StateMachine.Trigger.CLOSE, sessionId);
doCoolingDownAndClose(sessionId, "bridge_detach");
}
} catch (StateMachine.TransitionException e) {
log.warn("Bridge-triggered close transition failed: {}", e.getMessage());
/** mcphub.control.bridge_attach — increment bridge count; suspend idle timer on first attach. */
private JsonNode handleBridgeAttach(JsonNode params) {
long pid = params != null && params.has("pid") ? params.get("pid").asLong(-1) : -1;
// Deduplicate: same PID re-attaching is idempotent
boolean isNew = pid >= 0 && bridgePids.add(pid);
if (isNew) {
int count = activeBridgeCount.incrementAndGet();
bridgeLastPing.put(pid, Instant.now());
if (count == 1) {
sessionManager.suspendIdleTimer();
}
ObjectNode r = mapper.createObjectNode();
r.put("state", stateMachine.getState().name());
r.put("mcphub_providers", "stop");
return r;
log.info("Bridge attached (pid={}). Active bridges: {}", pid, count);
} else if (pid >= 0) {
// Re-attach from same PID: update ping timestamp only
bridgeLastPing.put(pid, Instant.now());
log.debug("Bridge re-attached (pid={}, idempotent). Active bridges: {}", pid, activeBridgeCount.get());
}
ObjectNode r = mapper.createObjectNode();
r.put("status", "ok");
r.put("active_bridges", activeBridgeCount.get());
return r;
}

/** mcphub.control.bridge_ping — update last-ping timestamp for a known bridge. */
private JsonNode handleBridgePing(JsonNode params) {
long pid = params != null && params.has("pid") ? params.get("pid").asLong(-1) : -1;
if (pid >= 0 && bridgePids.contains(pid)) {
bridgeLastPing.put(pid, Instant.now());
}
ObjectNode r = mapper.createObjectNode();
r.put("status", "ok");
return r;
}

/** mcphub.control.bridge_detach — decrement bridge count; resume idle timer when last bridge leaves. */
private JsonNode handleBridgeDetach(JsonNode params) {
long pid = params != null && params.has("pid") ? params.get("pid").asLong(-1) : -1;
// Only decrement if this PID was a known attached bridge
boolean wasKnown = pid >= 0 && bridgePids.remove(pid);
if (wasKnown) {
bridgeLastPing.remove(pid);
int count = activeBridgeCount.updateAndGet(c -> c > 0 ? c - 1 : 0);
log.info("Bridge detached (pid={}). Active bridges: {}", pid, count);
if (count == 0) {
sessionManager.resumeIdleTimer();
bridgePids.clear();
bridgeLastPing.clear();
}
} else {
log.debug("Bridge detach for unknown pid={}, ignoring", pid);
}
ObjectNode r = mapper.createObjectNode();
r.put("status", "ok");
r.put("active_bridges", activeBridgeCount.get());
return r;
}

// -------------------------------------------------------------------------
// Helpers
// -------------------------------------------------------------------------
Expand Down Expand Up @@ -444,4 +504,31 @@ private void doCoolingDownAndClose(String sessionId, String trigger) {
log.warn("CoolingDown→Closed transition failed: {}", e.getMessage());
}
}

/** Janitor: detect crashed bridges and clean up counts. */
private void janitorTask() {
try {
Instant now = Instant.now();
for (Long pid : new ArrayList<>(bridgePids)) {
Instant lastPing = bridgeLastPing.get(pid);
if (lastPing == null) {
lastPing = Instant.now();
}
if (now.getEpochSecond() - lastPing.getEpochSecond() > 180) {
boolean alive = ProcessHandle.of(pid).map(ProcessHandle::isAlive).orElse(false);
if (!alive) {
log.warn("Bridge pid {} appears dead, removing", pid);
bridgePids.remove(pid);
bridgeLastPing.remove(pid);
int count = activeBridgeCount.updateAndGet(c -> c > 0 ? c - 1 : 0);
if (count == 0) {
sessionManager.resumeIdleTimer();
}
}
}
}
} catch (Exception e) {
log.warn("Bridge janitor task failed: {}", e.getMessage());
}
}
}
23 changes: 21 additions & 2 deletions java/src/main/java/dev/sorted/mcphub/SessionManager.java
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ public interface TimeoutCallback {
private TimeoutCallback timeoutCallback;
private ScheduledFuture<?> idleTimerFuture;
private ScheduledFuture<?> armTimerFuture;
private volatile boolean idleSuspended;

public SessionManager() {
this(DEFAULT_IDLE_TIMEOUT_SECONDS, DEFAULT_ARM_TIMEOUT_SECONDS);
Expand All @@ -62,6 +63,7 @@ public String startSession() {
currentSessionId = UUID.randomUUID().toString();
sessionStart = Instant.now();
lastActivity = sessionStart;
idleSuspended = false;
log.info("Session started: {}", currentSessionId);
startArmTimer();
return currentSessionId;
Expand All @@ -79,6 +81,19 @@ public void resetActivity() {
resetIdleTimer();
}

/** Suspend idle timer while a bridge is attached. */
public synchronized void suspendIdleTimer() {
cancelIdleTimer();
idleSuspended = true;
}

/** Resume idle timer after last bridge detaches. */
public synchronized void resumeIdleTimer() {
idleSuspended = false;
lastActivity = Instant.now();
resetIdleTimer();
}

/** Increment in-flight call count. */
public int incrementInFlight() { return inFlightCount.incrementAndGet(); }

Expand All @@ -92,6 +107,7 @@ public void resetActivity() {
public void endSession() {
cancelIdleTimer();
cancelArmTimer();
idleSuspended = false;
String sid = currentSessionId;
currentSessionId = null;
sessionStart = null;
Expand Down Expand Up @@ -127,7 +143,10 @@ private void cancelArmTimer() {
if (armTimerFuture != null) { armTimerFuture.cancel(false); armTimerFuture = null; }
}

private void resetIdleTimer() {
private synchronized void resetIdleTimer() {
if (idleSuspended) {
return;
}
cancelIdleTimer();
idleTimerFuture = scheduler.schedule(() -> {
String sid = currentSessionId;
Expand All @@ -138,7 +157,7 @@ private void resetIdleTimer() {
}, idleTimeoutSeconds, TimeUnit.SECONDS);
}

private void cancelIdleTimer() {
private synchronized void cancelIdleTimer() {
if (idleTimerFuture != null) { idleTimerFuture.cancel(false); idleTimerFuture = null; }
}
}
72 changes: 71 additions & 1 deletion java/src/main/java/dev/sorted/mcphub/StdioBridge.java
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
public class StdioBridge {
private static final Logger log = LoggerFactory.getLogger(StdioBridge.class);
private static final ObjectMapper mapper = new ObjectMapper();
private volatile boolean attached = false;

/**
* Run the bridge loop. Blocks until stdin is closed.
Expand All @@ -40,6 +41,11 @@ public void run() throws IOException {

log.info("mcphub bridge started");

// Start heartbeat thread (best-effort 60s ping)
Thread heartbeat = new Thread(this::heartbeatLoop, "mcphub-bridge-heartbeat");
heartbeat.setDaemon(true);
heartbeat.start();

String line;
while ((line = stdinReader.readLine()) != null) {
if (line.isBlank()) continue;
Expand All @@ -64,8 +70,18 @@ public void run() throws IOException {
writeError(stdoutWriter, id, -32000,
"MCPHUB daemon is not running. Start it with 'mcphub start'.");
}

// Send bridge_attach on first successful initialize
if (!attached) {
JsonNode methodNode = req.get("method");
if (methodNode != null && "initialize".equals(methodNode.asText())) {
attached = sendAttach();
}
}
}

heartbeat.interrupt();

// REQ-2.3.5: stdin closed → bridge detaching (session stays open)
log.info("stdin closed, bridge detaching");
sendDetach();
Expand Down Expand Up @@ -102,13 +118,67 @@ private String forwardToDaemon(String requestLine) {
}
}

/** Send bridge_attach notification to daemon. Returns true on success. */
private boolean sendAttach() {
ObjectNode req = mapper.createObjectNode();
req.put("jsonrpc", "2.0");
req.put("method", "mcphub.control.bridge_attach");
ObjectNode params = mapper.createObjectNode();
params.put("pid", ProcessHandle.current().pid());
req.set("params", params);
try {
String response = forwardToDaemon(mapper.writeValueAsString(req));
if (response == null) {
log.warn("bridge_attach failed: daemon unreachable");
return false;
}
log.info("Bridge attached to daemon");
return true;
} catch (Exception e) {
log.warn("bridge_attach failed: {}", e.getMessage());
return false;
}
}

/** Send periodic bridge_ping to daemon (best-effort). */
private void sendPing() {
ObjectNode req = mapper.createObjectNode();
req.put("jsonrpc", "2.0");
req.put("method", "mcphub.control.bridge_ping");
ObjectNode params = mapper.createObjectNode();
params.put("pid", ProcessHandle.current().pid());
req.set("params", params);
try {
String response = forwardToDaemon(mapper.writeValueAsString(req));
if (response == null) {
log.warn("bridge_ping notification failed: daemon unreachable");
}
} catch (Exception e) {
log.warn("bridge_ping notification failed: {}", e.getMessage());
}
}

private void heartbeatLoop() {
while (!Thread.currentThread().isInterrupted()) {
try {
Thread.sleep(60000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
sendPing();
}
}

/** Send session-detach notification to daemon (best-effort). */
private void sendDetach() {
// Notify daemon so it can track last-bridge-exit and auto-close.
ObjectNode req = mapper.createObjectNode();
req.put("jsonrpc", "2.0");
req.put("method", "mcphub.control.bridge_detach");
req.set("params", mapper.createObjectNode());
ObjectNode params = mapper.createObjectNode();
params.put("pid", ProcessHandle.current().pid());
req.set("params", params);
try {
String response = forwardToDaemon(mapper.writeValueAsString(req));
if (response == null) {
Expand Down
Loading
Loading