mxtymoshyk opened a new pull request, #40393:
URL: https://github.com/apache/beam/pull/40393

   Fixes #18723.
   
   Fixes a race in two `MqttIOTest` read tests and re-enables 
`testReadNoClientId`, which has been `@Ignore`d under #18723. Test-only change.
   
   **Root cause**
   
   Before publishing, both tests called `doConnect(connection -> 
!connection.getConnectionId().isEmpty())`, which waits until the broker has any 
client with a non-empty connection id. The test's own publish client connects 
before the publisher thread starts, so the condition was already true and the 
wait returned after one 1-second poll. If the pipeline's reader had not 
subscribed by then, the broker had no subscriber for the topic and dropped the 
messages.
   
   - `testReadWithMetadata` reads for at most 10 seconds and then fails with 
every record missing. That is the failure on a recent `Java_IOs_Direct` run:
     ```
     java.lang.AssertionError: 
MqttIO.Read/Read(UnboundedMqttSource)/StripIds/ParMultiDo(StripIds).output:
     Expected: iterable with items [<MqttRecord{topic=topic/1, ...}>, ...]
          but: no item matches: <MqttRecord{topic=topic/1, ...}>, ...
     ```
     #34133 added a 2-second sleep for this (#34175). The sleep made the race 
less likely without removing it.
   - `testReadNoClientId` has no max read time, so the reader waits forever for 
records that never arrive and JUnit kills it at 60 seconds. That matches the 
timeout reported in #18723.
   
   **Fix**
   
   A new helper, `waitForSubscription`, polls the broker's topic subscriptions 
every 100 ms until one matches the destination the test is about to publish to. 
`Subscription.matches` handles wildcards, so the `topic/#` filter in 
`testReadWithMetadata` matches `topic.1`. The publish clients never subscribe, 
so a match means the reader has subscribed. The helper gives up after 30 
seconds.
   
   Both tests use it in place of `doConnect`, and the 2-second sleep is gone. 
The publisher threads now log exceptions instead of swallowing them, so a 
failed wait shows up in the test log. `testRead` already waits for the reader's 
own client id and is unchanged.
   
   **Repro**
   
   Add a delay to the reader in `MqttIO.UnboundedMqttReader.start()`, between 
`client.connect()` and the subscribe call, to stand in for a slow runner:
   
   ```java
   client.connect();
   Thread.sleep(5000);
   ```
   
   | | Before | After |
   |---|---|---|
   | `testReadWithMetadata`, reader delayed 5s | fails with `no item matches` | 
passes |
   | `testReadNoClientId`, reader delayed 5s | `test timed out after 60000 
milliseconds` | passes |
   | `testReadWithMetadata`, no delay, 20 runs | | 20/20, about 1.4s each (4.2s 
before) |
   | `testReadNoClientId`, no delay, 20 runs | | 20/20, about 1.4s each |
   
   Locally, `spotlessCheck`, `checkstyleTest` and `:sdks:java:io:mqtt:test` 
pass: 21 tests, none skipped.
   
   ------------------------
   
   Thank you for your contribution! Follow this checklist to help us 
incorporate your contribution quickly and easily:
   
    - [x] Mention the appropriate issue in your description (for example: 
`addresses #123`), if applicable. This will automatically add a link to the 
pull request in the issue. If you would like the issue to automatically close 
on merging the pull request, comment `fixes #<ISSUE NUMBER>` instead.
    - [ ] Update `CHANGES.md` with noteworthy changes.
    - [ ] If this contribution is large, please file an Apache [Individual 
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).
   
   See the [Contributor Guide](https://beam.apache.org/contribute) for more 
tips on [how to make review process 
smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier).
   


-- 
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]

Reply via email to