lymerin opened a new issue, #7417: URL: https://github.com/apache/shenyu/issues/7417
### Is there an existing issue for this? - [x] I have searched the existing issues ### Current Behavior `DividePluginCases.testRocketMQHello` sometimes ends with `Failed to consume RocketMQ access log` after its 30-second Awaitility wait. A successful HTTP response does not establish that the test's RocketMQ consumer has acquired the queues for `shenyu-access-logging`. On the pinned source commit `a12aa3e595e59f2b82ea2e73835770fe1fdcb88a` (which includes #7343 and #7227), I reproduced a specific version of this failure with the unmodified E2E test. The gateway returned HTTP 200, collected and sent the access log, and the broker stored that exact message. At the deadline the test consumer still had no target-topic queue assignment. The assertion reported a missing log although the message was in the broker. ### Expected Behavior Before sending the HTTP request, the test should confirm that its consumer owns the target topic queues. If it fails, the report should distinguish HTTP failure, broker storage failure, and consumer readiness/consumption failure. ### Steps To Reproduce Requirements: Linux/WSL2, JDK 17, Docker Compose v2, Python 3 with PyYAML, `nsenter`, `tc`, and Docker permissions. Save the three helper files from the collapsible source section below under `$AUDIT_ROOT`; no attachment is required. Use a fresh broker for each trial. 1. Pin the original source and build the images: ```bash export AUDIT_ROOT="$HOME/shenyu-mq-repro" export M2_REPO="$HOME/.m2/repository" mkdir -p "$AUDIT_ROOT" git clone https://github.com/apache/shenyu.git "$AUDIT_ROOT/source-linux" cd "$AUDIT_ROOT/source-linux" git checkout a12aa3e595e59f2b82ea2e73835770fe1fdcb88a git rev-parse HEAD > "$AUDIT_ROOT/SOURCE_SHA" ./mvnw -B clean install -Prelease,docker -Dmaven.javadoc.skip=true \ -Drat.skip=true -Dmaven.test.skip=true -Djacoco.skip=true \ -DskipITs -DskipTests -T2 ./mvnw -B -f shenyu-examples/shenyu-examples-http/pom.xml \ clean package -Pexample -DskipTests -Dmaven.javadoc.skip=true docker pull mysql:8.0 docker pull zookeeper:3.9 docker pull rocketmqinc/rocketmq:4.4.0 ``` 2. Prepare database/Compose files outside the checkout: ```bash mkdir -p "$AUDIT_ROOT/staging/shenyu-e2e/mysql/"{schema,driver} cp db/init/mysql/schema.sql "$AUDIT_ROOT/staging/shenyu-e2e/mysql/schema/schema.sql" printf "\nGRANT ALL PRIVILEGES ON shenyu.* TO 'shenyue2e'@'%%';\n" \ >> "$AUDIT_ROOT/staging/shenyu-e2e/mysql/schema/schema.sql" curl -fL --retry 3 \ https://repo.maven.apache.org/maven2/mysql/mysql-connector-java/8.0.29/mysql-connector-java-8.0.29.jar \ -o "$AUDIT_ROOT/staging/shenyu-e2e/mysql/driver/mysql-connector-java-8.0.29.jar" python3 "$AUDIT_ROOT/prepare-compose.py" "$AUDIT_ROOT" ``` 3. Run a control, compile the independent broker reader, then run the timing condition: ```bash sudo -E env AUDIT_ROOT="$AUDIT_ROOT" M2_REPO="$M2_REPO" \ TEST_FILTER=DividePluginTest,LoggingRuleSyncTest \ bash "$AUDIT_ROOT/run-baseline.sh" http control ``` The control creates the Surefire report needed for the test classpath. Compile `MqDump.java` once: ```bash mkdir -p "$AUDIT_ROOT/dump" python3 - "$AUDIT_ROOT" <<'PY' import pathlib, sys, xml.etree.ElementTree as ET root = pathlib.Path(sys.argv[1]) report = root / "source-linux/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/target/surefire-reports/TEST-org.apache.shenyu.e2e.testcase.logging.rocketmq.DividePluginTest.xml" properties = {item.attrib["name"]: item.attrib["value"] for item in ET.parse(report).getroot().findall("properties/property")} (root / "dump/classpath.txt").write_text(properties["java.class.path"]) PY javac --release 17 -cp "$(cat "$AUDIT_ROOT/dump/classpath.txt")" \ -d "$AUDIT_ROOT/dump" "$AUDIT_ROOT/MqDump.java" ``` Run the delayed trial: ```bash sudo -E env AUDIT_ROOT="$AUDIT_ROOT" M2_REPO="$M2_REPO" \ TEST_FILTER=DividePluginTest,LoggingRuleSyncTest \ FAULT=backend-delay BACKEND_DELAY_MS=1200 \ bash "$AUDIT_ROOT/run-baseline.sh" http backend-repro ``` The delay affects only the HTTP example container’s outgoing response. Java test code, MQ network, rebalance settings, and the 30-second wait stay unchanged. Replace `http` with `websocket` or `zookeeper` to test those sync modes. Logs are saved under `$AUDIT_ROOT/results/`; the harness removes its trial containers and volumes. `TEST_FILTER` excludes an unrelated test that calls an external website. 4. Check `maven.log` for the 30-second `testRocketMQHello` timeout, `backend-delay-counters.txt` for transmitted packets without intentional loss, and `broker-messages-java17.txt` for `/http/order/findById?id=23` with `status=200`. The `MqDump.java` reader uses an independent group and explicit offsets, so it does not change the test consumer’s progress. Image/setup or HTTP assertion failures are different outcomes. Repeat with a fresh broker if the timing window is missed. <details><summary>Tested helper scripts (copy into <code>$AUDIT_ROOT</code>)</summary> **`prepare-compose.py`** ```python #!/usr/bin/env python3 """Diagnostic deployment configuration. Does not edit the source checkout.""" import pathlib import sys import yaml root = pathlib.Path(sys.argv[1] if len(sys.argv) > 1 else "/opt/shenyu-ci-audit") source = root / "source-linux" cases = source / "shenyu-e2e/shenyu-e2e-case" case = cases / "shenyu-e2e-case-logging-rocketmq/compose" output = root / "compose" output.mkdir(exist_ok=True) def environment(service, extra): original = service.get("environment", {}) if isinstance(original, list): original = dict(item.split("=", 1) for item in original) service["environment"] = {**original, **extra} for mode in ["websocket", "http", "zookeeper"]: result = {"services": {}, "networks": {"shenyu": {"name": "shenyu", "external": True}}} for file in [cases / f"compose/sync/shenyu-sync-{mode}.yml", case / "shenyu-rocketmq-compose.yml", case / "shenyu-examples-http-compose.yml"]: data = yaml.safe_load(file.read_text()) result["services"].update(data["services"]) for name, service in result["services"].items(): service["volumes"] = [volume.replace("/tmp/shenyu-e2e", str(root / "staging/shenyu-e2e")) for volume in service.get("volumes", [])] service["restart"] = "no" service["ports"] = [f"127.0.0.1:{port}" for port in service.get("ports", [])] service.pop("deploy", None) service["mem_limit"] = "768m" service["memswap_limit"] = "1536m" if name == "shenyu-admin": environment(service, {"ADMIN_JVM": "-Xms64m -Xmx384m -XX:MaxDirectMemorySize=64m -XX:ActiveProcessorCount=2"}) elif name == "shenyu-bootstrap": environment(service, {"BOOT_JVM": "-Xms64m -Xmx384m -XX:MaxDirectMemorySize=96m -XX:ActiveProcessorCount=2"}) elif name == "shenyu-examples-http": environment(service, {"JAVA_TOOL_OPTIONS": "-Xms32m -Xmx160m -XX:MaxDirectMemorySize=32m -XX:ActiveProcessorCount=2"}) elif name == "rocketmq-dialevoneid": environment(service, {"JAVA_OPT_EXT": "-Xms64m -Xmx128m -Xmn32m"}) elif name == "rocketmq-broker": environment(service, {"JAVA_OPT_EXT": "-Xms64m -Xmx192m -Xmn64m"}) elif name == "shenyu-zookeeper": environment(service, {"JVMFLAGS": "-Xms32m -Xmx96m"}) elif name == "shenyu-mysql": service["command"] = ["--innodb-buffer-pool-size=64M", "--performance-schema=OFF", "--max-connections=80"] (output / f"{mode}.yml").write_text(yaml.safe_dump(result, sort_keys=False)) print(output / f"{mode}.yml") ``` **`run-baseline.sh`** ```bash #!/usr/bin/env bash # Repeated unmodified RocketMQ E2E tests, with external resource overrides. set -Eeuo pipefail ROOT=${AUDIT_ROOT:-/opt/shenyu-ci-audit} SOURCE=$ROOT/source-linux MODE=${1:-http} LABEL=${2:-baseline} case "$MODE" in http|websocket|zookeeper) ;; *) exit 2;; esac RUN=$ROOT/results/$(date -u +%Y%m%dT%H%M%SZ)-$MODE-$LABEL mkdir -p "$RUN" COMPOSE=${AUDIT_COMPOSE:-$ROOT/compose/$MODE.yml} DC=(docker compose -p shenyu-mq-audit -f "$COMPOSE") exec > >(tee "$RUN/harness.log") 2>&1 printf 'RUN=%s\nSHA=%s\nMODE=%s\n' "$RUN" "$(cat "$ROOT/SOURCE_SHA")" "$MODE" cp "$COMPOSE" "$RUN/compose.yml" date -u --iso-8601=ns uname -a free -h docker image inspect apache/shenyu-admin:latest apache/shenyu-bootstrap:latest rocketmqinc/rocketmq:4.4.0 > "$RUN/images.json" snapshot() { set +e date -u --iso-8601=ns > "$RUN/end-time.txt" docker stats --no-stream > "$RUN/docker-stats.txt" if [[ ${FAULT:-none} == backend-delay ]]; then example_pid=$(docker inspect --format '{{.State.Pid}}' shenyu-examples-http) nsenter -t "$example_pid" -n tc -s qdisc show dev eth0 > "$RUN/backend-delay-counters.txt" 2>&1 fi docker inspect $("${DC[@]}" ps -q -a) > "$RUN/containers.json" "${DC[@]}" logs --no-color --timestamps > "$RUN/containers.log" 2>&1 free -h > "$RUN/memory-end.txt" dmesg -T | tail -100 > "$RUN/kernel-tail.txt" for container in shenyu-bootstrap rocketmq-broker rocketmq-dialevoneid; do docker exec "$container" sh -c 'find /root /home/shenyu /home/rocketmq /opt -type f -name "*rocketmq*log*" 2>/dev/null' > "$RUN/$container-log-paths.txt" mkdir -p "$RUN/$container-files" docker cp "$container:/root/logs/rocketmqlogs" "$RUN/$container-files/" 2>/dev/null docker cp "$container:/home/shenyu/logs/rocketmqlogs" "$RUN/$container-files/" 2>/dev/null docker cp "$container:/home/rocketmq/logs/rocketmqlogs" "$RUN/$container-files/" 2>/dev/null done for endpoint in plugins selectorData ruleData; do curl -sS --max-time 5 "http://localhost:31195/actuator/$endpoint" > "$RUN/gateway-$endpoint.json" done mkdir -p "$RUN/test-reports" "$RUN/client-logs" cp -a /root/logs/rocketmqlogs/. "$RUN/client-logs/" 2>/dev/null find "$SOURCE/shenyu-e2e" -path '*/target/surefire-reports/*' -type f -exec cp --parents '{}' "$RUN/test-reports/" \; timeout 25 docker exec -e JAVA_OPT_EXT='-Xms32m -Xmx96m -Xmn32m' rocketmq-broker sh mqadmin topicStatus -n rocketmq-dialevoneid:9876 -t shenyu-access-logging > "$RUN/topic-status.txt" 2>&1 timeout 25 docker exec -e JAVA_OPT_EXT='-Xms32m -Xmx96m -Xmn32m' rocketmq-broker sh mqadmin consumerProgress -n rocketmq-dialevoneid:9876 -g shenyu-plugin-logging-rocketmq > "$RUN/consumer-progress.txt" 2>&1 timeout 30 docker exec -e JAVA_OPT_EXT='-Xms32m -Xmx96m -Xmn32m' rocketmq-broker sh mqadmin printMsg -n rocketmq-dialevoneid:9876 -t shenyu-access-logging -d true -c UTF-8 > "$RUN/broker-messages.txt" 2>&1 timeout 30 java -Xms32m -Xmx96m -cp "$ROOT/dump:$(cat "$ROOT/dump/classpath.txt")" MqDump > "$RUN/broker-messages-java17.txt" 2>&1 docker cp rocketmq-broker:/home/rocketmq/store/config "$RUN/broker-store-config" 2>/dev/null sha256sum "$SOURCE/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/src/test/java/org/apache/shenyu/e2e/testcase/logging/rocketmq/DividePluginCases.java" > "$RUN/test-source-sha256.txt" if [[ ${KEEP_ENV:-0} != 1 ]]; then "${DC[@]}" down -v --remove-orphans > "$RUN/cleanup.log" 2>&1 fi } trap 'rc=$?; echo "$rc" > "$RUN/harness.exit"; snapshot; exit "$rc"' EXIT cd "$SOURCE" docker network inspect shenyu >/dev/null 2>&1 || docker network create shenyu "${DC[@]}" up -d --pull never shenyu-mysql shenyu-admin shenyu-bootstrap sh shenyu-e2e/shenyu-e2e-case/k8s/script/healthcheck.sh http://localhost:31095/actuator/health sh shenyu-e2e/shenyu-e2e-case/k8s/script/healthcheck.sh http://localhost:31195/actuator/health "${DC[@]}" up -d --pull never rocketmq-dialevoneid rocketmq-broker for port in 31876 10911; do ready=0 for attempt in $(seq 1 30); do if (echo > /dev/tcp/localhost/$port) 2>/dev/null; then ready=1; break; fi sleep 2 done [[ $ready == 1 ]] || { echo "TCP readiness failed: $port"; exit 3; } done "${DC[@]}" up -d --pull never shenyu-examples-http sh shenyu-e2e/shenyu-e2e-case/k8s/script/healthcheck.sh http://localhost:31189/actuator/health sleep 10 date -u --iso-8601=ns > "$RUN/test-start-time.txt" if [[ ${FAULT:-none} == backend-delay ]]; then EXAMPLE_PID=$(docker inspect --format '{{.State.Pid}}' shenyu-examples-http) nsenter -t "$EXAMPLE_PID" -n tc qdisc replace dev eth0 root netem delay "${BACKEND_DELAY_MS:-1200}ms" nsenter -t "$EXAMPLE_PID" -n tc -s qdisc show dev eth0 > "$RUN/backend-delay-config.txt" printf 'Injected only backend egress delay: %sms; test and MQ traffic unchanged\n' "${BACKEND_DELAY_MS:-1200}" fi export MAVEN_OPTS='-Xms32m -Xmx256m -XX:ActiveProcessorCount=2' TEST_ARGS=() if [[ -n ${M2_REPO:-} ]]; then TEST_ARGS+=("-Dmaven.repo.local=$M2_REPO") fi if [[ -n ${TEST_FILTER:-} ]]; then TEST_ARGS+=("-Dtest=$TEST_FILTER" -Dsurefire.failIfNoSpecifiedTests=false) fi ARG_LINE='-Xms32m -Xmx256m -XX:ActiveProcessorCount=2' if [[ ${TRACE:-0} == 1 ]]; then ARG_LINE+=" -javaagent:$ROOT/trace/trace-agent.jar" fi FAULT_PID='' if [[ ${FAULT:-none} == nameserver ]]; then python3 "$ROOT/inject-nameserver-fault.py" --run "$RUN" --duration "${FAULT_DURATION:-2}" > "$RUN/fault-watch.log" 2>&1 & FAULT_PID=$! fi set +e ./mvnw -B -f ./shenyu-e2e/pom.xml -pl shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq -am test "-DargLine=$ARG_LINE" "${TEST_ARGS[@]}" > "$RUN/maven.log" 2>&1 RC=$? set -e echo "$RC" > "$RUN/maven.exit" if [[ -n "$FAULT_PID" ]]; then wait "$FAULT_PID" || echo 'Fault watch did not complete successfully'; fi tail -65 "$RUN/maven.log" exit "$RC" ``` **`MqDump.java`** ```java import java.nio.charset.StandardCharsets; import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer; import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.message.MessageQueue; /** Post-test observation only; uses an independent group and explicit offsets. */ public class MqDump { public static void main(String[] args) throws Exception { DefaultMQPullConsumer consumer = new DefaultMQPullConsumer("shenyu-ci-audit-dump"); consumer.setNamesrvAddr("localhost:31876"); consumer.start(); try { for (MessageQueue queue : consumer.fetchSubscribeMessageQueues("shenyu-access-logging")) { long min = consumer.minOffset(queue); long max = consumer.maxOffset(queue); System.out.println("QUEUE " + queue + " minOffset=" + min + " maxOffset=" + max); long cursor = min; while (cursor < max) { PullResult result = consumer.pull(queue, "*", cursor, 32); System.out.println("PULL " + queue + " " + result.getPullStatus() + " next=" + result.getNextBeginOffset()); if (result.getMsgFoundList() != null) { for (MessageExt message : result.getMsgFoundList()) { System.out.println("MESSAGE id=" + message.getMsgId() + " queue=" + message.getQueueId() + " offset=" + message.getQueueOffset() + " bornEpochMs=" + message.getBornTimestamp() + " storeEpochMs=" + message.getStoreTimestamp() + " body=" + new String(message.getBody(), StandardCharsets.UTF_8)); } } if (result.getNextBeginOffset() <= cursor) break; cursor = result.getNextBeginOffset(); } } } finally { consumer.shutdown(); } } } ``` </details> ### Environment ```markdown - ShenYu `2.7.2-SNAPSHOT`, pinned commit `a12aa3e595e59f2b82ea2e73835770fe1fdcb88a`. - Ubuntu 24.04 under WSL2, Java 17, Docker 29.1.3 / Compose v2.40.3 in the executed trials. - RocketMQ broker/nameserver image `rocketmqinc/rocketmq:4.4.0`; test/producer client 4.9.3. - The generated Compose file uses smaller JVM heaps and locally bound ports for WSL. It does not alter the Java test, MQ network, or rebalance interval. ``` ### Debug logs With a 1200 ms delay on the HTTP example’s outgoing response, the **unchanged** RocketMQ E2E assertion timed out in all three sync modes. An independent post-test broker read found the exact HTTP 200 access log in each trial: | Sync mode | MQ test | `upstreamResponseTime` | Broker log | Delayed packets / dropped | | --- | --- | ---: | --- | ---: | | HTTP | 30 s timeout | 1213 ms | URI and status 200 found | 125 / 0 | | WebSocket | 30 s timeout | 1215 ms | URI and status 200 found | 94 / 0 | | ZooKeeper | 30 s timeout | 1213 ms | URI and status 200 found | 94 / 0 | All three Maven logs contain: ```text org.awaitility.core.ConditionTimeoutException: ... within 30 seconds. [ERROR] DividePluginTest.testDivide:146->lambda$testDivide$2:146 Failed to consume RocketMQ access log ``` `websocket/broker-messages-java17.txt` independently reports `queueId=2`, `minOffset=0`, `maxOffset=1`, `PULL ... FOUND next=1`, and a message at `offset=0` with `requestUri=http://localhost:31195/http/order/findById?id=23`, `status=200`, `upstreamResponseTime=1215`. The delay counter says `delay 1.2s` and `94 pkt (dropped 0)`. The traced WebSocket trial gives this sequence relative to `consumer.start()` returning: ```text +0.015 s initial rebalance: no access-log topic route or target queue +0.744 s broker notification: target topic still absent +1.361 s broker stores the exact HTTP 200 access log (queue 2, offset 0) +20.769 s target route is learned; processQueueTable still has only the retry-topic queue original 30 s wait expires; test shuts down its consumer ``` The gateway trace records `collect` → `consume0` → `sendOneway`. The trial ends when the test shuts down, so it does **not** claim a later consumption at 40 seconds. A separate review of 2026-09-01 through 2026-10-02 E2E logs found 93 failed RocketMQ jobs among 1,156 completed jobs, affecting 36 PRs. That is a **mixed job failure rate**, not the rate of this queue race: it includes rule-sync, setup, and infrastructure failures. Nineteen historical records show a matched logging rule followed by no consumed log (16 using the 30-second await; three with older sleep-based code). Historical CI logs lack the broker/queue snapshots collected here. ### Anything else? **Why this happens:** the test starts a default push consumer before the first access log creates `shenyu-access-logging`. In RocketMQ client 4.9.3, [the clustering rebalance path](https://github.com/apache/rocketmq/blob/rocketmq-all-4.9.3/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java#L266) reads its cached topic queue set before `findConsumerIdList` can fetch missing route information. The iteration that learns the route can still use the previously empty queue set. The [default rebalance interval](https://github.com/apache/rocketmq/blob/rocketmq-all-4.9.3/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceService.java#L28) is 20 seconds. In the reproduced ordering, the next ordinary queue assignment is beyond the test’s 30-second wait. `consumer.start()` is therefore not a sufficient readiness check. Related work: [#6817](https://github.com/apache/shenyu/issues/6817) described fixed sleeps in the earlier E2E test; [#6978](https://github.com/apache/shenyu/pull/6978) replaced them with Awaitility; [#7343](https://github.com/apache/shenyu/pull/7343) added logging-rule synchronization; [#7227](https://github.com/apache/shenyu/pull/7227) made collector startup idempotent. The pinned reproduction commit includes these changes. This issue concerns the remaining test-consumer preparation path. **Historical CI examples with the same 30-second consumption timeout:** | PR / failing job | Mode | Observation | What it does not establish | | --- | --- | --- | --- | | [#7263, job 110208034111](https://github.com/apache/shenyu/actions/runs/36811415051/job/110208034111) | HTTP | Logging selector/rule matched, then the MQ assertion timed out (~31.2 s). Checkout included the #7343 rule wait. | #7227 was not yet present; no broker/queue dump. | | [#7104, job 110213639462](https://github.com/apache/shenyu/actions/runs/36813219874/job/110213639462) | ZooKeeper | Logging rule matched, then the MQ assertion timed out (~30.7 s). Checkout included #7343. | Same evidence limit. | | [#7343, job 110160238387](https://github.com/apache/shenyu/actions/runs/36795833090/job/110160238387) | ZooKeeper | Logging rule matched on this PR branch, then the MQ assertion timed out (~31.0 s). | Branch was still being developed; no broker/queue state. | | [#7236, job 110133574738](https://github.com/apache/shenyu/actions/runs/36677759249/job/110133574738) | HTTP | Logging rule matched, then the MQ assertion timed out (~31.0 s). | Older checkout lacked both #7343 and #7227. | These links show that the symptom affected unrelated PRs. They do **not** prove that each historical timeout had the same broker/consumer state. The 30-day inventory also includes rule-sync, environment-preparation, and infrastructure failures, so its overall job failure rate is not the rate of this particular race. [#7289 also had an earlier RocketMQ timeout](https://github.com/apache/shenyu/actions/runs/36415934554/job/108915531794), but that run predates #7343 and its rule arrived late; it is **not** evidence for this queue race. The decisive evidence for this issue is the controlled run above: an unchanged test timed out while an independent broker read recovered its exact HTTP 200 log and client logs showed no target queue assignment. No external evidence archive is required to understand that result. -- 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]
