YvCeung commented on code in PR #7903:
URL: https://github.com/apache/incubator-seata/pull/7903#discussion_r2777296972
##########
server/src/main/java/org/apache/seata/server/cluster/manager/ClusterWatcherManager.java:
##########
@@ -128,31 +248,131 @@ private void sendWatcherResponse(Watcher<HttpContext>
watcher, HttpResponseStatu
} else {
ctx.writeAndFlush(response);
}
- } else {
- // HTTP/2 response (h2c support)
- // Send headers frame
+ return;
+ }
+
+ // For HTTP/2, headers must be sent first on the initial response
+ if (sendHeaders) {
Http2Headers headers = new
DefaultHttp2Headers().status(nettyStatus.codeAsText());
- headers.set(HttpHeaderNames.CONTENT_LENGTH, "0");
+ headers.set(HttpHeaderNames.CONTENT_TYPE, "text/event-stream;
charset=utf-8");
+ headers.set(HttpHeaderNames.CACHE_CONTROL, "no-cache");
+
ctx.write(new DefaultHttp2HeadersFrame(headers));
+ }
- // Send empty data frame with endStream=true to close the stream
- ctx.writeAndFlush(new DefaultHttp2DataFrame(Unpooled.EMPTY_BUFFER,
true))
- .addListener(f -> {
- if (!f.isSuccess()) {
- logger.warn("HTTP2 response send failed,
group={}", group, f.cause());
- }
- });
+ String group = watcher.getGroup();
+ String eventData = buildEventData(group);
+ ByteBuf content = Unpooled.copiedBuffer(eventData,
StandardCharsets.UTF_8);
+
+ // Send DATA frame (if closeStream is true, it will end the current
stream)
+ ctx.write(new DefaultHttp2DataFrame(content, closeStream));
+ ctx.flush();
+ }
+
+ /**
+ * Get current cluster metadata for the given group.
+ * This method extracts the logic from ClusterController#cluster to avoid
circular dependency.
+ *
+ * @param group the group name
+ * @return the MetadataResponse containing current cluster metadata
+ */
+ public MetadataResponse getMetadataResponse(String group) {
+ MetadataResponse metadataResponse = new MetadataResponse();
+ if (StringUtils.isBlank(group)) {
+ group = ConfigurationFactory.getInstance()
+ .getConfig(ConfigurationKeys.SERVER_RAFT_GROUP,
DEFAULT_SEATA_GROUP);
+ }
+ RaftServer raftServer = RaftServerManager.getRaftServer(group);
+ if (raftServer != null) {
+ String mode =
ConfigurationFactory.getInstance().getConfig(STORE_MODE);
+ metadataResponse.setStoreMode(mode);
+ RouteTable routeTable = RouteTable.getInstance();
+ try {
+
routeTable.refreshLeader(RaftServerManager.getCliClientServiceInstance(),
group, 1000);
+ PeerId leader = routeTable.selectLeader(group);
+ if (leader != null) {
+ Set<Node> nodes = new HashSet<>();
+ RaftClusterMetadata raftClusterMetadata =
+
raftServer.getRaftStateMachine().getRaftLeaderMetadata();
+ Node leaderNode = raftServer
+ .getRaftStateMachine()
+ .getRaftLeaderMetadata()
+ .getLeader();
+ leaderNode.setGroup(group);
+ nodes.add(leaderNode);
+ nodes.addAll(raftClusterMetadata.getLearner());
+ nodes.addAll(raftClusterMetadata.getFollowers());
+ metadataResponse.setTerm(raftClusterMetadata.getTerm());
+ metadataResponse.setNodes(new ArrayList<>(nodes));
+ }
+ } catch (Exception e) {
+ logger.error("Failed to get cluster metadata for group {}:
{}", group, e.getMessage(), e);
+ }
+ }
+ return metadataResponse;
+ }
+
+ /**
+ * Build event data with simplified format: only group, timestamp, and
metadata.
+ * For HTTP/2 connections, this will send the full MetadataResponse.
+ *
+ * @param group the group name
+ * @return the event data string with prefix
+ */
+ private String buildEventData(String group) {
+ try {
+ // Get current MetadataResponse
+ MetadataResponse metadataResponse = getMetadataResponse(group);
+
+ // Build simplified JSON: only group, timestamp, and metadata
+ String json = String.format(
+ "{\"group\":\"%s\",\"timestamp\":%d,\"metadata\":%s}",
+ group,
+ System.currentTimeMillis(),
+ OBJECT_MAPPER.writeValueAsString(metadataResponse));
+
+ logger.debug("Sending watch event: group={}, term={}", group,
metadataResponse.getTerm());
+ return Constants.WATCH_EVENT_PREFIX + json + "\n";
+ } catch (JsonProcessingException e) {
+ logger.error("Failed to serialize MetadataResponse for group {}:
{}", group, e.getMessage(), e);
+ // Fallback: send minimal data
+ String json = String.format(
+ "{\"group\":\"%s\",\"timestamp\":%d,\"metadata\":null}",
+ group, System.currentTimeMillis());
+ return Constants.WATCH_EVENT_PREFIX + json + "\n";
}
}
public void registryWatcher(Watcher<HttpContext> watcher) {
String group = watcher.getGroup();
Long term = GROUP_UPDATE_TERM.get(group);
- if (term == null || watcher.getTerm() >= term) {
- WATCHERS.computeIfAbsent(group, value -> new
ConcurrentLinkedQueue<>())
+ HttpContext context = watcher.getAsyncContext();
+ boolean isHttp2 = context.isHttp2();
+
+ // For HTTP/2, must send response headers immediately, cannot delay
+ if (isHttp2 && !HTTP2_HEADERS_SENT.getOrDefault(watcher, false)) {
+ sendWatcherResponse(watcher, HttpResponseStatus.OK, false, true);
+ HTTP2_HEADERS_SENT.put(watcher, true);
+ }
+
+ if (isHttp2) {
+ HTTP2_WATCHERS
+ .computeIfAbsent(group, value -> new
ConcurrentLinkedQueue<>())
.add(watcher);
+
+ // If term has been updated, notify immediately
+ if (term != null && term > watcher.getTerm()) {
+ notifyWatcher(watcher, null);
Review Comment:
在当前实现下不会出现该竞态,原因有两点:
执行顺序:registryWatcher 里是先执行 **发 headers + HTTP2_HEADERS_SENT.put(watcher,
true)**(338–341 行),再把 **watcher 加入 HTTP2_WATCHERS(**343–346 行),之后才可能在同一调用内或由
onChangeEvent 触发 notifyWatcher。
调用关系:能执行到 notifyWatcher 的 watcher,要么来自同一次 registryWatcher**(此时本线程已执行过
put)**,要么来自 onChangeEvent 遍历队列(此时该 watcher 已在队中,说明入队前已做过 put)。因此两处看到的都是「已发过
headers」,不会重复发送。
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]