diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClient.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClient.java index 62834ecb6dde..4d45b517d086 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClient.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClient.java @@ -321,9 +321,9 @@ private Mono acquireSpecificOrNextSession(String * *

The returned {@link PagedFlux} fetches additional pages from the broker on demand using * cursor-based pagination (server-returned {@code skip} plus {@code lastSessionId} of the - * previous page) and terminates when the broker returns an empty page. The default page size - * is 100; callers can request a different size via - * {@link PagedFlux#byPage(int)}.

+ * previous page) and terminates when the broker returns a page smaller than the requested page + * size (a short or empty page signals the end). The default page size is 100; callers can + * request a different size via {@link PagedFlux#byPage(int)}.

* * @return A {@link PagedFlux} of session ID strings. */ @@ -339,9 +339,9 @@ public PagedFlux listSessions() { * *

The returned {@link PagedFlux} fetches additional pages from the broker on demand using * cursor-based pagination (server-returned {@code skip} plus {@code lastSessionId} of the - * previous page) and terminates when the broker returns an empty page. The default page size - * is 100; callers can request a different size via - * {@link PagedFlux#byPage(int)}.

+ * previous page) and terminates when the broker returns a page smaller than the requested page + * size (a short or empty page signals the end). The default page size is 100; callers can + * request a different size via {@link PagedFlux#byPage(int)}.

* *

Values at or beyond the active-messages sentinel value * ({@code new Date(253402300800000L)}, rendered by {@code OffsetDateTime.toString()} as @@ -440,13 +440,15 @@ private Mono> fetchSessionPage(OffsetDateTime lastUpdatedT managementNode -> managementNode.getMessageSessions(lastUpdatedTime, skip, pageSize, lastSessionId)) .map(result -> { final java.util.List sessionIds = result.getSessionIds(); - // Empty page terminates pagination (matches Track 1's SessionBrowser loop and the - // broker contract). Continuation token encodes the server-returned skip and the - // last session ID of the page so the next call uses the same cursor Track 1 does. - // Base64url-encode the session ID so arbitrary byte sequences (including the '|' - // separator) round-trip without escaping. + // A short page (fewer IDs than the requested page size) terminates pagination: the + // service has no more sessions to return, so no further page is fetched. This + // matches the .NET SDK, which breaks on page.Count < SessionBrowsePageSize. + // Continuation token encodes the server-returned skip and the last session ID of the + // page so the next call uses the same cursor Track 1 does. Base64url-encode the + // session ID so arbitrary byte sequences (including the '|' separator) round-trip + // without escaping. final String continuationToken; - if (sessionIds.isEmpty()) { + if (sessionIds.size() < pageSize) { continuationToken = null; } else { final String last = sessionIds.get(sessionIds.size() - 1); diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusManagementNode.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusManagementNode.java index d706a085b011..e05facd030f0 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusManagementNode.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusManagementNode.java @@ -155,8 +155,8 @@ Mono updateDisposition(String lockToken, DispositionStatus dispositionStat *

Pagination follows the cursor semantics of Track 1's * {@code com.microsoft.azure.servicebus.SessionBrowser}: the caller threads {@code skip} from * {@link MessageSessionsResult#getNextSkip()} of the previous response and {@code lastSessionId} - * (the last entry of the previous page) into the next request, and stops when an empty page is - * returned.

+ * (the last entry of the previous page) into the next request, and stops when the broker returns + * a page smaller than the requested page size (a short or empty page signals the end).

* * @param lastUpdatedTime Filter timestamp. To get sessions with active messages, pass the * {@link ManagementConstants#ACTIVE_MESSAGES_SENTINEL} sentinel (the implementation also diff --git a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClientTest.java b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClientTest.java index 7b5cb95ce692..a535665e38ee 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClientTest.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/test/java/com/azure/messaging/servicebus/ServiceBusSessionReceiverAsyncClientTest.java @@ -43,10 +43,14 @@ import java.time.Duration; import java.time.OffsetDateTime; import java.time.ZoneOffset; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.List; import java.util.concurrent.Callable; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; +import java.util.stream.IntStream; import static com.azure.messaging.servicebus.ReceiverOptions.createNamedSessionOptions; import static com.azure.messaging.servicebus.ReceiverOptions.createUnnamedSessionOptions; @@ -383,23 +387,42 @@ private ServiceBusSessionReceiverAsyncClient newSessionReceiver() { }, CLIENT_IDENTIFIER, false); } + /** + * Builds a full page of {@code count} session IDs ("{prefix}0" .. "{prefix}{count-1}") as a + * mutable list. A full page (count == the requested page size) drives a further page request + * under short-page pagination termination. + */ + private static List fullPage(String prefix, int count) { + return IntStream.range(0, count).mapToObj(i -> prefix + i).collect(Collectors.toCollection(ArrayList::new)); + } + + private static List fullPage(String prefix) { + return fullPage(prefix, 100); + } + /** * Verifies the no-arg listSessions() drives the broker with the active-messages sentinel and - * collects every page until the broker returns an empty page. + * collects every page until the broker returns a short page (fewer IDs than the requested page + * size), which terminates pagination. */ @Test - void listSessionsActiveModeStreamsAllPagesUntilEmpty() { - // First page: 2 sessions, server-returned skip = 2. + void listSessionsActiveModeStreamsAllPagesUntilShortPage() { + // First page: a full page (100 sessions) continues; server-returned skip = 100. + final List firstPage = fullPage("s"); when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(0), eq(100), - isNull())).thenReturn(Mono.just(new MessageSessionsResult(Arrays.asList("s1", "s2"), 2))); - // Cursor for the second page is encoded server-skip + base64url(lastSessionId), which decodes - // back to (skip=2, lastSessionId="s2"). Empty page terminates pagination. - when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(2), eq(100), - eq("s2"))).thenReturn(Mono.just(new MessageSessionsResult(Collections.emptyList(), 2))); + isNull())).thenReturn(Mono.just(new MessageSessionsResult(firstPage, 100))); + // Cursor for the second page is server-skip (100) + base64url(lastSessionId "s99"). The second + // page is short (2 < 100), which terminates pagination. + when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(100), eq(100), + eq("s99"))).thenReturn(Mono.just(new MessageSessionsResult(Arrays.asList("t1", "t2"), 102))); final ServiceBusSessionReceiverAsyncClient client = newSessionReceiver(); - StepVerifier.create(client.listSessions()).expectNext("s1", "s2").expectComplete().verify(DEFAULT_TIMEOUT); + StepVerifier.create(client.listSessions()) + .expectNextSequence(firstPage) + .expectNext("t1", "t2") + .expectComplete() + .verify(DEFAULT_TIMEOUT); } /** @@ -411,20 +434,22 @@ void listSessionsActiveModeStreamsAllPagesUntilEmpty() { void listSessionsHonorsServerSkipAndLastSessionId() { final OffsetDateTime sessionStateUpdatedAfter = OffsetDateTime.of(2026, 1, 1, 0, 0, 0, 0, ZoneOffset.UTC); - // First page returns 2 items but the server reports skip = 7 (5 entries filtered server-side). + // First page returns a full page of 100 items but the server reports skip = 107 (extra entries + // filtered server-side), so the second-page request must use the server-returned skip, not + // requestSkip + page.size(). + final List firstPage = fullPage("a"); when(managementNode.getMessageSessions(eq(sessionStateUpdatedAfter), eq(0), eq(100), isNull())) - .thenReturn(Mono.just(new MessageSessionsResult(Arrays.asList("a", "b"), 7))); - // Second-page request must use the server-returned skip (7) and lastSessionId ("b"). - when(managementNode.getMessageSessions(eq(sessionStateUpdatedAfter), eq(7), eq(100), eq("b"))) - .thenReturn(Mono.just(new MessageSessionsResult(Collections.singletonList("c"), 8))); - // Third page empty terminates pagination. - when(managementNode.getMessageSessions(eq(sessionStateUpdatedAfter), eq(8), eq(100), eq("c"))) - .thenReturn(Mono.just(new MessageSessionsResult(Collections.emptyList(), 8))); + .thenReturn(Mono.just(new MessageSessionsResult(firstPage, 107))); + // Second-page request must use the server-returned skip (107) and lastSessionId ("a99"). It is + // short (1 < 100), which terminates pagination. + when(managementNode.getMessageSessions(eq(sessionStateUpdatedAfter), eq(107), eq(100), eq("a99"))) + .thenReturn(Mono.just(new MessageSessionsResult(Collections.singletonList("c"), 108))); final ServiceBusSessionReceiverAsyncClient client = newSessionReceiver(); StepVerifier.create(client.listSessions(sessionStateUpdatedAfter)) - .expectNext("a", "b", "c") + .expectNextSequence(firstPage) + .expectNext("c") .expectComplete() .verify(DEFAULT_TIMEOUT); } @@ -450,14 +475,23 @@ void listSessionsRejectsNullSessionStateUpdatedAfter() { void listSessionsRoundTripsArbitrarySessionIdsThroughCursor() { final String sessionWithPipe = "weird|session|id"; + // First page is full (100 items) with the pipe-containing id LAST, so the cursor for the next + // page encodes it; a full page also drives the second request under short-page termination. + final List firstPage = fullPage("x", 99); + firstPage.add(sessionWithPipe); when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(0), eq(100), - isNull())).thenReturn(Mono.just(new MessageSessionsResult(Collections.singletonList(sessionWithPipe), 1))); - when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(1), eq(100), - eq(sessionWithPipe))).thenReturn(Mono.just(new MessageSessionsResult(Collections.emptyList(), 1))); + isNull())).thenReturn(Mono.just(new MessageSessionsResult(firstPage, 100))); + // The second-page request must decode the cursor back to lastSessionId=sessionWithPipe intact + // (pipe and all); the short (empty) page then terminates pagination. + when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(100), eq(100), + eq(sessionWithPipe))).thenReturn(Mono.just(new MessageSessionsResult(Collections.emptyList(), 100))); final ServiceBusSessionReceiverAsyncClient client = newSessionReceiver(); - StepVerifier.create(client.listSessions()).expectNext(sessionWithPipe).expectComplete().verify(DEFAULT_TIMEOUT); + StepVerifier.create(client.listSessions()) + .expectNextSequence(firstPage) + .expectComplete() + .verify(DEFAULT_TIMEOUT); } /** @@ -525,22 +559,24 @@ void listSessionsEmptyContinuationTokenCompletes() { */ @Test void listSessionsHonorsCallerPageSize() { - // Caller asks for pages of 25; both first-page and next-page calls must pass top=25. + // Caller asks for pages of 25; both first-page and next-page calls must pass top=25. The first + // page is full (25 items) so a second page is requested; the short second page (1 < 25) + // terminates pagination. + final List firstPage = fullPage("s", 25); when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(0), eq(25), - isNull())).thenReturn(Mono.just(new MessageSessionsResult(Arrays.asList("a", "b"), 2))); - when( - managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(2), eq(25), eq("b"))) - .thenReturn(Mono.just(new MessageSessionsResult(Collections.emptyList(), 2))); + isNull())).thenReturn(Mono.just(new MessageSessionsResult(firstPage, 25))); + when(managementNode.getMessageSessions(eq(ManagementConstants.ACTIVE_MESSAGES_SENTINEL), eq(25), eq(25), + eq("s24"))).thenReturn(Mono.just(new MessageSessionsResult(Collections.singletonList("t1"), 26))); final ServiceBusSessionReceiverAsyncClient client = newSessionReceiver(); StepVerifier.create(client.listSessions().byPage(25)).assertNext(page -> { - org.junit.jupiter.api.Assertions.assertEquals(2, page.getValue().size()); - org.junit.jupiter.api.Assertions.assertEquals("a", page.getValue().get(0)); - org.junit.jupiter.api.Assertions.assertEquals("b", page.getValue().get(1)); - }) - .assertNext(page -> org.junit.jupiter.api.Assertions.assertTrue(page.getValue().isEmpty())) - .expectComplete() - .verify(DEFAULT_TIMEOUT); + org.junit.jupiter.api.Assertions.assertEquals(25, page.getValue().size()); + org.junit.jupiter.api.Assertions.assertEquals("s0", page.getValue().get(0)); + org.junit.jupiter.api.Assertions.assertEquals("s24", page.getValue().get(24)); + }).assertNext(page -> { + org.junit.jupiter.api.Assertions.assertEquals(1, page.getValue().size()); + org.junit.jupiter.api.Assertions.assertEquals("t1", page.getValue().get(0)); + }).expectComplete().verify(DEFAULT_TIMEOUT); } }