Skip to content
Open
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
18 changes: 17 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,10 +107,26 @@ This endpoint requires authentication (standard Jenkins permissions) and provide
"streamableRequestsTotal": 150,
"connectionErrorsTotal": 2,
"uptimeSeconds": 3600,
"startTime": "2025-01-28T10:00:00Z"
"startTime": "2025-01-28T10:00:00Z",
"users": {
"alice": {
"sseConnectionsTotal": 10,
"sseConnectionsActive": 1,
"streamableRequestsTotal": 50,
"connectionErrorsTotal": 0
},
"bob": {
"sseConnectionsTotal": 32,
"sseConnectionsActive": 2,
"streamableRequestsTotal": 100,
"connectionErrorsTotal": 2
}
}
}
```

The top-level counters are server-wide aggregates. The `users` object breaks down the same counters per Jenkins user (including `anonymous` and `SYSTEM`), allowing you to see which accounts are actively connected and how much traffic each is generating.

#### Graceful Shutdown

When Jenkins shuts down, the health endpoint will return `503 Service Unavailable` with a brief grace period before full termination. This allows clients to detect the shutdown and prepare for reconnection.
Expand Down
13 changes: 8 additions & 5 deletions src/main/java/io/jenkins/plugins/mcp/server/Endpoint.java
Original file line number Diff line number Diff line change
Expand Up @@ -277,7 +277,8 @@
response.sendError(HttpServletResponse.SC_NOT_FOUND, "Streamable endpoint is disabled");
return true;
}
McpConnectionMetrics.recordStreamableRequest();
McpConnectionMetrics.recordStreamableRequest(
Jenkins.getAuthentication2().getName());
handleMessage(request, response, httpServletStreamableServerTransportProvider);
return true;
}
Expand Down Expand Up @@ -537,7 +538,8 @@
if (isBrowserRequest(req)) {
serveBrowserPage(resp);
} else {
McpConnectionMetrics.recordStreamableRequest();
McpConnectionMetrics.recordStreamableRequest(
Jenkins.getAuthentication2().getName());
handleMessage(req, resp, httpServletStreamableServerTransportProvider);
}
return true;
Expand Down Expand Up @@ -718,16 +720,17 @@
private void handleSSE(HttpServletRequest request, HttpServletResponse response)
throws IOException, ServletException {
String clientInfo = getClientInfo(request);
String username = Jenkins.getAuthentication2().getName();
log.info("SSE connection started from {}", clientInfo);
McpConnectionMetrics.recordSseConnectionStart();
McpConnectionMetrics.recordSseConnectionStart(username);
try {
httpServletSseServerTransportProvider.service(request, response);
} catch (IOException e) {
log.warn("SSE connection error from {}: {}", clientInfo, e.getMessage());
McpConnectionMetrics.recordConnectionError();
McpConnectionMetrics.recordConnectionError(username);

Check warning on line 730 in src/main/java/io/jenkins/plugins/mcp/server/Endpoint.java

View check run for this annotation

ci.jenkins.io / Code Coverage

Not covered line

Line 730 is not covered by tests
throw e;
} finally {
McpConnectionMetrics.recordSseConnectionEnd();
McpConnectionMetrics.recordSseConnectionEnd(username);
log.info("SSE connection ended from {}", clientInfo);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import java.io.IOException;
import java.time.Duration;
import java.time.Instant;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import jenkins.model.Jenkins;
import lombok.extern.slf4j.Slf4j;
Expand All @@ -55,7 +56,15 @@
* "streamableRequestsTotal": 150,
* "connectionErrorsTotal": 2,
* "uptimeSeconds": 3600,
* "startTime": "2025-01-28T10:00:00Z"
* "startTime": "2025-01-28T10:00:00Z",
* "users": {
* "alice": {
* "sseConnectionsTotal": 10,
* "sseConnectionsActive": 1,
* "streamableRequestsTotal": 50,
* "connectionErrorsTotal": 0
* }
* }
* }
* </pre>
*/
Expand Down Expand Up @@ -92,6 +101,22 @@ public class McpConnectionMetrics {
*/
private static final Instant startTime = Instant.now();

/**
* Per-user metrics, keyed by Jenkins username (including "anonymous" and "SYSTEM").
*/
static final ConcurrentHashMap<String, UserMetrics> userMetrics = new ConcurrentHashMap<>();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

what happen for instance with large number of users? (possible memory leak?)

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hi @olamy
Thank you for taking a look at the PR. You're right that entries are never removed and the map grows once per unique username. Since this is bounded by the number of Jenkins users rather than by traffic, I think the memory footprint should stay small in practice (~150 bytes per user). I can add a comment to that effect if you'd like. Thoughts?


private static UserMetrics userMetricsFor(String username) {
return userMetrics.computeIfAbsent(username, k -> new UserMetrics());
}

static final class UserMetrics {
final AtomicLong sseConnectionsTotal = new AtomicLong(0);
final AtomicLong sseConnectionsActive = new AtomicLong(0);
final AtomicLong streamableRequestsTotal = new AtomicLong(0);
final AtomicLong connectionErrorsTotal = new AtomicLong(0);
}

/**
* Handles GET requests to the metrics endpoint.
* Requires authentication (checks Jenkins.READ permission).
Expand Down Expand Up @@ -128,6 +153,17 @@ public static void handleMetricsRequest(HttpServletResponse response) throws IOE
responseJson.put("uptimeSeconds", uptime.getSeconds());
responseJson.put("startTime", startTime.toString());

ObjectNode usersNode = objectMapper.createObjectNode();
userMetrics.forEach((username, metrics) -> {
ObjectNode userNode = objectMapper.createObjectNode();
userNode.put("sseConnectionsTotal", metrics.sseConnectionsTotal.get());
userNode.put("sseConnectionsActive", metrics.sseConnectionsActive.get());
userNode.put("streamableRequestsTotal", metrics.streamableRequestsTotal.get());
userNode.put("connectionErrorsTotal", metrics.connectionErrorsTotal.get());
usersNode.set(username, userNode);
});
responseJson.set("users", usersNode);

response.getWriter().write(objectMapper.writeValueAsString(responseJson));
response.getWriter().flush();
}
Expand All @@ -140,27 +176,61 @@ public static void recordSseConnectionStart() {
sseConnectionsActive.incrementAndGet();
}

/**
* Records a new SSE connection starting for the given user.
*/
public static void recordSseConnectionStart(String username) {
recordSseConnectionStart();
UserMetrics u = userMetricsFor(username);
u.sseConnectionsTotal.incrementAndGet();
u.sseConnectionsActive.incrementAndGet();
}

/**
* Records an SSE connection ending.
*/
public static void recordSseConnectionEnd() {
sseConnectionsActive.decrementAndGet();
}

/**
* Records an SSE connection ending for the given user.
*/
public static void recordSseConnectionEnd(String username) {
recordSseConnectionEnd();
userMetricsFor(username).sseConnectionsActive.decrementAndGet();
}

/**
* Records a Streamable HTTP request.
*/
public static void recordStreamableRequest() {
streamableRequestsTotal.incrementAndGet();
}

/**
* Records a Streamable HTTP request for the given user.
*/
public static void recordStreamableRequest(String username) {
recordStreamableRequest();
userMetricsFor(username).streamableRequestsTotal.incrementAndGet();
}

/**
* Records a connection error.
*/
public static void recordConnectionError() {
connectionErrorsTotal.incrementAndGet();
}

/**
* Records a connection error for the given user.
*/
public static void recordConnectionError(String username) {
recordConnectionError();
userMetricsFor(username).connectionErrorsTotal.incrementAndGet();
}

/**
* Returns the current number of active SSE connections.
*
Expand Down Expand Up @@ -205,5 +275,6 @@ static void reset() {
sseConnectionsActive.set(0);
streamableRequestsTotal.set(0);
connectionErrorsTotal.set(0);
userMetrics.clear();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -164,4 +164,106 @@ void testReset() {
assertThat(McpConnectionMetrics.getTotalStreamableRequests()).isZero();
assertThat(McpConnectionMetrics.getTotalConnectionErrors()).isZero();
}

@Test
void testPerUserSseConnectionMetrics() {
McpConnectionMetrics.recordSseConnectionStart("alice");
McpConnectionMetrics.recordSseConnectionStart("alice");
McpConnectionMetrics.recordSseConnectionStart("bob");

McpConnectionMetrics.UserMetrics alice = McpConnectionMetrics.userMetrics.get("alice");
McpConnectionMetrics.UserMetrics bob = McpConnectionMetrics.userMetrics.get("bob");

assertThat(alice).isNotNull();
assertThat(alice.sseConnectionsTotal.get()).isEqualTo(2);
assertThat(alice.sseConnectionsActive.get()).isEqualTo(2);

assertThat(bob).isNotNull();
assertThat(bob.sseConnectionsTotal.get()).isOne();
assertThat(bob.sseConnectionsActive.get()).isOne();

// Global counters should also reflect the combined total
assertThat(McpConnectionMetrics.getTotalSseConnections()).isEqualTo(3);
assertThat(McpConnectionMetrics.getActiveSseConnections()).isEqualTo(3);

McpConnectionMetrics.recordSseConnectionEnd("alice");
assertThat(alice.sseConnectionsActive.get()).isOne();
assertThat(McpConnectionMetrics.getActiveSseConnections()).isEqualTo(2);
}

@Test
void testPerUserStreamableRequestMetrics() {
McpConnectionMetrics.recordStreamableRequest("carol");
McpConnectionMetrics.recordStreamableRequest("carol");
McpConnectionMetrics.recordStreamableRequest("dave");

McpConnectionMetrics.UserMetrics carol = McpConnectionMetrics.userMetrics.get("carol");
McpConnectionMetrics.UserMetrics dave = McpConnectionMetrics.userMetrics.get("dave");

assertThat(carol).isNotNull();
assertThat(carol.streamableRequestsTotal.get()).isEqualTo(2);

assertThat(dave).isNotNull();
assertThat(dave.streamableRequestsTotal.get()).isOne();

assertThat(McpConnectionMetrics.getTotalStreamableRequests()).isEqualTo(3);
}

@Test
void testPerUserConnectionErrorMetrics() {
McpConnectionMetrics.recordConnectionError("alice");
McpConnectionMetrics.recordConnectionError("alice");
McpConnectionMetrics.recordConnectionError("bob");

McpConnectionMetrics.UserMetrics alice = McpConnectionMetrics.userMetrics.get("alice");
McpConnectionMetrics.UserMetrics bob = McpConnectionMetrics.userMetrics.get("bob");

assertThat(alice).isNotNull();
assertThat(alice.connectionErrorsTotal.get()).isEqualTo(2);

assertThat(bob).isNotNull();
assertThat(bob.connectionErrorsTotal.get()).isOne();

assertThat(McpConnectionMetrics.getTotalConnectionErrors()).isEqualTo(3);
}

@Test
void testResetClearsUserMetrics() {
McpConnectionMetrics.recordSseConnectionStart("alice");
McpConnectionMetrics.recordStreamableRequest("bob");

assertThat(McpConnectionMetrics.userMetrics).isNotEmpty();

McpConnectionMetrics.reset();

assertThat(McpConnectionMetrics.userMetrics).isEmpty();
}

@Test
void testMetricsEndpointIncludesUsersSection(JenkinsRule jenkins) throws Exception {
McpConnectionMetrics.recordSseConnectionStart("alice");
McpConnectionMetrics.recordStreamableRequest("alice");
McpConnectionMetrics.recordStreamableRequest("bob");

try (JenkinsRule.WebClient webClient = jenkins.createWebClient()) {
var metricsUrl = jenkins.getURL().toString() + McpConnectionMetrics.URL_NAME;
var request = new WebRequest(new URL(metricsUrl), HttpMethod.GET);
WebResponse response = webClient.loadWebResponse(request);

assertThat(response.getStatusCode()).isEqualTo(HttpServletResponse.SC_OK);

DocumentContext json = JsonPath.parse(response.getContentAsString());
assertThat(json.read("$.users", Object.class)).isNotNull();
assertThat(json.read("$.users.alice.sseConnectionsTotal", Long.class))
.isEqualTo(1L);
assertThat(json.read("$.users.alice.sseConnectionsActive", Long.class))
.isEqualTo(1L);
assertThat(json.read("$.users.alice.streamableRequestsTotal", Long.class))
.isEqualTo(1L);
assertThat(json.read("$.users.alice.connectionErrorsTotal", Long.class))
.isZero();
assertThat(json.read("$.users.bob.streamableRequestsTotal", Long.class))
.isEqualTo(1L);
}
}
}
Loading