diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java index 01fc25185f7e..cecfd6a9b626 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/ResumableUploadFutureImpl.java @@ -247,9 +247,44 @@ private void fail(Throwable t) { if (inFlight != null) { inFlight.cancel(true); } - progressTracker.onFailed(t, uploadSessionUrl); + Throwable augmented = augmentWithUrl(t); + progressTracker.onFailed(augmented, uploadSessionUrl); closePayload(); - resultFuture.setException(t); + resultFuture.setException(augmented); + } + + private Throwable augmentWithUrl(Throwable t) { + String url = uploadSessionUrl; + if (url == null || url.isEmpty()) { + return t; + } + String message = t.getMessage(); + if (message != null && message.contains(url)) { + return t; + } + String baseMessage = message != null ? message : t.getClass().getSimpleName(); + String augmentedMessage = baseMessage + " (upload URL: " + url + ")"; + Throwable augmented = t; + if (t instanceof ApiException) { + ApiException apiException = (ApiException) t; + augmented = + ApiExceptionFactory.createException( + augmentedMessage, + apiException, + apiException.getStatusCode(), + apiException.isRetryable(), + apiException.getErrorDetails()); + } else if (t instanceof IllegalStateException) { + augmented = new IllegalStateException(augmentedMessage, t); + } else if (t instanceof IOException) { + augmented = new IOException(augmentedMessage, t); + } + if (augmented != t) { + for (Throwable suppressed : t.getSuppressed()) { + augmented.addSuppressed(suppressed); + } + } + return augmented; } private void closePayload() { diff --git a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/RewindableStreamBuffer.java b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/RewindableStreamBuffer.java index 24a27dafabcf..73ce3aef0a79 100644 --- a/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/RewindableStreamBuffer.java +++ b/sdk-platform-java/gax-java/gax/src/main/java/com/google/api/gax/rpc/RewindableStreamBuffer.java @@ -101,7 +101,8 @@ void realignTo(long committedOffset) throws IOException { throw UploadErrors.protocolViolation( String.format( "Server committed offset %d is below buffer base offset %d for upload URL %s; cannot" - + " rewind stream before buffer base", + + " rewind stream before buffer base. A seekable stream is required to rewind to" + + " earlier offsets.", committedOffset, bufferBaseOffset, uploadUrl)); } diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java index 86f5512fa088..74d47bdf506c 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/ResumableUploadCallableImplTest.java @@ -110,7 +110,11 @@ void setUp() { callContext = FakeCallContext.createDefault(); executor = Executors.newScheduledThreadPool(2); clientContext = - ClientContext.newBuilder().setDefaultCallContext(callContext).setExecutor(executor).build(); + ClientContext.newBuilder() + .setDefaultCallContext(callContext) + .setExecutor(executor) + .setEndpoint("https://test.endpoint.com") + .build(); callable = new ResumableUploadCallableImpl<>(mockClient, defaultSettings, clientContext); } @@ -1504,6 +1508,138 @@ void testProgressListener_orderingUnderConcurrency_pinsSequentialExecutor() thro } } + @Test + void testActionableErrors_startFailure_preservesOriginalExceptionWithoutEndpointSuffix() { + ApiException startError = createApiException(401, StatusCode.Code.UNAUTHENTICATED); + when(mockStartCallable.futureCall(any(), any())) + .thenReturn(ApiFutures.immediateFailedFuture(startError)); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isSameInstanceAs(startError); + assertThat(ex.getCause().getMessage()).doesNotContain("endpoint:"); + assertThat(future.getUploadSessionUrl()).isNull(); + } + + @Test + void testActionableErrors_chunkFailure_messageContainsUploadSessionUrl() { + String sessionUrl = "https://upload.url/chunk-error-test"; + stubStartSession(sessionUrl); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(403, StatusCode.Code.PERMISSION_DENIED))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isInstanceOf(ApiException.class); + assertThat(ex.getCause().getMessage()).contains(sessionUrl); + assertThat(future.getUploadSessionUrl()).isEqualTo(sessionUrl); + } + + @Test + void testActionableErrors_preservesErrorDetailsCauseChainAndSuppressedExceptions() { + String sessionUrl = "https://upload.url/chunk-error-details-test"; + stubStartSession(sessionUrl); + ErrorDetails errorDetails = ErrorDetails.builder().build(); + ApiException original = + ApiExceptionFactory.createException( + "HTTP 403", + null, + new HttpStatusStatusCode(403, StatusCode.Code.PERMISSION_DENIED), + false, + errorDetails); + IOException suppressed = new IOException("underlying stream error"); + original.addSuppressed(suppressed); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn(ApiFutures.immediateFailedFuture(original)); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isInstanceOf(ApiException.class); + ApiException cause = (ApiException) ex.getCause(); + assertThat(cause.getMessage()).contains(sessionUrl); + assertThat(cause.getCause()).isSameInstanceAs(original); + assertThat(cause.getErrorDetails()).isSameInstanceAs(errorDetails); + assertThat(cause.getSuppressed()).asList().contains(suppressed); + assertThat(future.getUploadSessionUrl()).isEqualTo(sessionUrl); + } + + @Test + void testActionableErrors_recoveryFailure_messageContainsUploadSessionUrl() { + String sessionUrl = "https://upload.url/recovery-error-test"; + stubStartSession(sessionUrl); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(400, StatusCode.Code.INVALID_ARGUMENT))); + when(mockQueryCallable.futureCall(any(QueryStatusRequest.class), any())) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(403, StatusCode.Code.PERMISSION_DENIED))); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isInstanceOf(ApiException.class); + assertThat(ex.getCause().getMessage()).contains(sessionUrl); + assertThat(future.getUploadSessionUrl()).isEqualTo(sessionUrl); + } + + @Test + void testActionableErrors_globalTimeoutFailure_messageContainsUploadSessionUrl() { + String sessionUrl = "https://upload.url/timeout-error-test"; + stubStartSession(sessionUrl); + SettableApiFuture> hungChunk = SettableApiFuture.create(); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())).thenReturn(hungChunk); + + ResumableUploadCallSettings settings = + defaultSettings.toBuilder().setGlobalTimeout(Duration.ofMillis(50)).build(); + + ResumableUploadFuture future = + callable.futureCall("resource-path", streamOf("hello"), null, settings); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isInstanceOf(DeadlineExceededException.class); + assertThat(ex.getCause().getMessage()).contains(sessionUrl); + assertThat(future.getUploadSessionUrl()).isEqualTo(sessionUrl); + } + + @Test + void testActionableErrors_rewindFailure_surfacesActionableSeekableStreamMessage() { + String sessionUrl = "https://upload.url/rewind-error-test"; + stubStartSession(sessionUrl); + when(mockChunkCallable.futureCall(any(ChunkUploadRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + ChunkUploadResponse.create(ResumableUploadStatus.ACTIVE, null))) + .thenReturn( + ApiFutures.immediateFailedFuture( + createApiException(400, StatusCode.Code.INVALID_ARGUMENT))); + + when(mockQueryCallable.futureCall(any(QueryStatusRequest.class), any())) + .thenReturn( + ApiFutures.immediateFuture( + createQueryResponse(false, 4L, null, ResumableUploadStatus.ACTIVE))); + + byte[] data = new byte[16]; + ResumableUploadFuture future = + callable.futureCall("resource-path", new ByteArrayInputStream(data), null); + + ExecutionException ex = assertThrows(ExecutionException.class, future::get); + assertThat(ex.getCause()).isInstanceOf(FailedPreconditionException.class); + assertThat(ex.getCause().getMessage()).contains(sessionUrl); + assertThat(ex.getCause().getMessage()).contains("seekable stream"); + assertThat(future.getUploadSessionUrl()).isEqualTo(sessionUrl); + } + private static class HttpStatusStatusCode implements StatusCode { private final int httpStatus; private final StatusCode.Code code; diff --git a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/RewindableStreamBufferTest.java b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/RewindableStreamBufferTest.java index 7b3ac1f7ce8a..df3bc3f25b4c 100644 --- a/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/RewindableStreamBufferTest.java +++ b/sdk-platform-java/gax-java/gax/src/test/java/com/google/api/gax/rpc/RewindableStreamBufferTest.java @@ -154,6 +154,7 @@ void testRealignToBelowBaseOffset_throwsFailedPreconditionException_classifiedFa assertThat(exception.getMessage()).contains("4"); assertThat(exception.getMessage()).contains("8"); assertThat(exception.getMessage()).contains(UPLOAD_URL); + assertThat(exception.getMessage()).contains("seekable stream"); // Must be classified as FATAL by ResumableUploadErrorClassifier ResumableUploadErrorClassifier.Category category =