From 48a2b905f4aa468f2ec54b9b20de403e65bb757f Mon Sep 17 00:00:00 2001 From: Ameya Shendre <52448570+ameya9@users.noreply.github.com> Date: Tue, 3 Mar 2026 16:13:44 +0000 Subject: [PATCH 01/11] docs: Add test for Session Invalidation (#1483) --- .../generic/GenericImporterTest.java | 29 ++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/extensions/data-transfer/portability-data-transfer-generic/src/test/java/org/datatransferproject/datatransfer/generic/GenericImporterTest.java b/extensions/data-transfer/portability-data-transfer-generic/src/test/java/org/datatransferproject/datatransfer/generic/GenericImporterTest.java index 1380b2ac8..af9ae5363 100644 --- a/extensions/data-transfer/portability-data-transfer-generic/src/test/java/org/datatransferproject/datatransfer/generic/GenericImporterTest.java +++ b/extensions/data-transfer/portability-data-transfer-generic/src/test/java/org/datatransferproject/datatransfer/generic/GenericImporterTest.java @@ -20,6 +20,7 @@ import org.datatransferproject.spi.cloud.storage.TemporaryPerJobDataStore; import org.datatransferproject.spi.transfer.idempotentexecutor.InMemoryIdempotentImportExecutor; import org.datatransferproject.spi.transfer.types.DestinationMemoryFullException; +import org.datatransferproject.spi.transfer.types.SessionInvalidatedException; import org.datatransferproject.transfer.JobMetadata; import org.datatransferproject.types.common.models.IdOnlyContainerResource; import org.datatransferproject.types.transfer.auth.AppCredentials; @@ -277,7 +278,7 @@ public void testGenericImporterBadRequest() throws Exception { assertEquals("itemId", error.title()); assertContains("(400) bad_request", error.exception()); } - + @Test public void testGenericImporterUnexpectedResponse() throws Exception { InMemoryIdempotentImportExecutor executor = new InMemoryIdempotentImportExecutor(monitor); @@ -426,4 +427,30 @@ public void testGenericImporterPassingJobMetadataRecurringJob() throws Exception assertTrue(executor.getErrors().isEmpty()); } + @Test + public void testSessionInvalidatedException() throws Exception { + InMemoryIdempotentImportExecutor executor = new InMemoryIdempotentImportExecutor(monitor); + GenericImporter importer = + getImporter( + importerClass, + container -> + List.of( + new ImportableData<>( + new GenericPayload<>(container.getId(), "schemasource"), + container.getId(), + container.getId()))); + webServer.enqueue(new MockResponse() + .setResponseCode(401) + .setBody("{\"error\":\"invalid_token\"}")); + webServer.enqueue(new MockResponse().setResponseCode(401).setBody("{\"error\":\"session_invalidated\"}")); + + assertThrows(SessionInvalidatedException.class, () -> { + importer.importItem( + UUID.randomUUID(), + executor, + new TokensAndUrlAuthData( + "accessToken", "refreshToken", webServer.url("/refresh").toString()), + new IdOnlyContainerResource("itemId")); + }); + } } From eb6362fbfdca0b5209448fced92df91397c27a5c Mon Sep 17 00:00:00 2001 From: "Chih-Hsien (Simon) Yeh" Date: Fri, 13 Mar 2026 10:45:15 +0800 Subject: [PATCH 02/11] feat: add SynologySignalHandler to Synology transfer extension (#1484) Implemented `SynologySignalHandler` to manage job lifecycle signals (e.g., job started or completion) within the Synology transfer extension. This ensures that the Synology C2 API is notified of the final state of a transfer job, allowing for better synchronization. ## Changes - New Signal Handler: Added `SynologySignalHandler` to process and retry job lifecycle signals using the RetryingCallable framework. - Service Integration: Added `sendJobSignal` to `SynologyDTPService` to handle the actual HTTP POST request to the Synology C2 API. - Extension Update: Modified `SynologyTransferExtension` to register the new SignalHandler during initialization. - Configuration & Models: - Updated `synology.yaml` to include the new `/import/job/signal` endpoint. - Updated `C2Api` and its `ApiPath` inner class to support the new signal path. - Testing: - Added `SynologySignalHandlerTest` to verify successful signal transmission and the retry mechanism on failure. - Updated `TestConfigs` to include the signal path for existing tests. --- .../synology/SynologyTransferExtension.java | 11 ++ .../datatransfer/synology/models/C2Api.java | 15 ++- .../synology/service/SynologyDTPService.java | 22 +++ .../signals/SynologySignalHandler.java | 112 +++++++++++++++ .../src/main/resources/config/synology.yaml | 1 + .../signals/SynologySignalHandlerTest.java | 127 ++++++++++++++++++ .../synology/utils/TestConfigs.java | 6 +- 7 files changed, 292 insertions(+), 2 deletions(-) create mode 100644 extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandler.java create mode 100644 extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/SynologyTransferExtension.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/SynologyTransferExtension.java index 94990986e..f5246eea8 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/SynologyTransferExtension.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/SynologyTransferExtension.java @@ -14,6 +14,7 @@ import org.datatransferproject.datatransfer.synology.photos.SynologyPhotosImporter; import org.datatransferproject.datatransfer.synology.service.SynologyDTPService; import org.datatransferproject.datatransfer.synology.service.SynologyOAuthTokenManager; +import org.datatransferproject.datatransfer.synology.signals.SynologySignalHandler; import org.datatransferproject.datatransfer.synology.uploader.SynologyUploader; import org.datatransferproject.datatransfer.synology.videos.SynologyVideosImporter; import org.datatransferproject.spi.cloud.storage.AppCredentialStore; @@ -23,6 +24,7 @@ import org.datatransferproject.spi.transfer.idempotentexecutor.IdempotentImportExecutorExtension; import org.datatransferproject.spi.transfer.provider.Exporter; import org.datatransferproject.spi.transfer.provider.Importer; +import org.datatransferproject.spi.transfer.provider.SignalHandler; import org.datatransferproject.transfer.JobMetadata; import org.datatransferproject.types.common.models.DataVertical; import org.datatransferproject.types.transfer.auth.AppCredentials; @@ -36,6 +38,7 @@ public class SynologyTransferExtension implements TransferExtension { private boolean initialized = false; private ImmutableMap importerMap; + private SynologySignalHandler signalHandler; @Override public String getServiceId() { @@ -64,6 +67,11 @@ private String getExtensionSecret() { return importerMap.get(transferDataType); } + @Override + public SignalHandler getSignalHandler() { + return signalHandler; + } + @Override public void initialize(ExtensionContext context) { if (initialized) { @@ -113,6 +121,9 @@ public void initialize(ExtensionContext context) { importerBuilder.put( VIDEOS, new SynologyVideosImporter(monitor, tokenManager, synologyUploader)); importerMap = importerBuilder.build(); + signalHandler = + new SynologySignalHandler( + tokenManager, synologyDTPService, context.getSetting("retryLibrary", null)); monitor.info(() -> "Initializing SynologyTransferExtension"); initialized = true; diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java index 1814abaf5..1cc3ef799 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java @@ -28,6 +28,7 @@ public class C2Api extends ServiceConfig.Service { private final String createAlbum; private final String uploadItem; private final String addItemToAlbum; + private final String signalJob; @JsonCreator public C2Api(@JsonProperty("baseUrl") String baseUrl, @JsonProperty("apiPath") ApiPath apiPath) { @@ -38,6 +39,7 @@ public C2Api(@JsonProperty("baseUrl") String baseUrl, @JsonProperty("apiPath") A this.createAlbum = UrlUtils.join(baseUrl, apiPath.getCreateAlbumPath()); this.uploadItem = UrlUtils.join(baseUrl, apiPath.getUploadItemPath()); this.addItemToAlbum = UrlUtils.join(baseUrl, apiPath.getAddItemToAlbumPath()); + this.signalJob = UrlUtils.join(baseUrl, apiPath.getSignalJobPath()); } public ApiPath getApiPath() { @@ -56,19 +58,26 @@ public String getAddItemToAlbum() { return addItemToAlbum; } + public String getSignalJob() { + return signalJob; + } + public static class ApiPath { private final String createAlbumPath; private final String uploadItemPath; private final String addItemToAlbumPath; + private final String signalJobPath; @JsonCreator public ApiPath( @JsonProperty("createAlbum") String createAlbumPath, @JsonProperty("uploadItem") String uploadItemPath, - @JsonProperty("addItemToAlbum") String addItemToAlbumPath) { + @JsonProperty("addItemToAlbum") String addItemToAlbumPath, + @JsonProperty("signalJob") String signalJobPath) { this.createAlbumPath = createAlbumPath; this.uploadItemPath = uploadItemPath; this.addItemToAlbumPath = addItemToAlbumPath; + this.signalJobPath = signalJobPath; } public String getCreateAlbumPath() { @@ -82,5 +91,9 @@ public String getUploadItemPath() { public String getAddItemToAlbumPath() { return addItemToAlbumPath; } + + public String getSignalJobPath() { + return signalJobPath; + } } } diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java index 038b0eff7..e946e2fda 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java @@ -44,6 +44,7 @@ import org.datatransferproject.spi.transfer.types.DestinationMemoryFullException; import org.datatransferproject.spi.transfer.types.NoNasInAccountException; import org.datatransferproject.spi.transfer.types.UploadErrorException; +import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle; import org.datatransferproject.types.common.models.media.MediaAlbum; import org.datatransferproject.types.common.models.photos.PhotoModel; import org.datatransferproject.types.common.models.videos.VideoModel; @@ -246,6 +247,27 @@ public Map addItemToAlbum(String albumId, String itemId, UUID jo return sendPostRequest(c2Api.getAddItemToAlbum(), builder.build(), jobId); } + /** + * Updates job status. + * + * @param jobStatus the job status + * @param jobId the job ID + * @return a map of shape {"success": } + */ + public Map sendJobSignal(JobLifeCycle jobStatus, UUID jobId) + throws CopyExceptionWithFailureReason { + FormBody.Builder builder = + new FormBody.Builder() + .add("job_id", jobId.toString()) + .add("service", exportingService) + .add("state", jobStatus.state().name()); + if (jobStatus.endReason() != null) { + builder.add("end_reason", jobStatus.endReason().name()); + } + + return sendPostRequest(c2Api.getSignalJob(), builder.build(), jobId); + } + @VisibleForTesting protected void throwExceptionIfNoQuota(Response response) throws CopyExceptionWithFailureReason { String errorCode = ""; diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandler.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandler.java new file mode 100644 index 000000000..67b4b314a --- /dev/null +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandler.java @@ -0,0 +1,112 @@ +/* + * Copyright 2026 The Data Transfer Project Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.datatransferproject.datatransfer.synology.signals; + +import java.time.Clock; +import java.util.Objects; +import java.util.UUID; +import java.util.concurrent.Callable; +import org.datatransferproject.api.launcher.Monitor; +import org.datatransferproject.datatransfer.synology.service.SynologyDTPService; +import org.datatransferproject.datatransfer.synology.service.SynologyOAuthTokenManager; +import org.datatransferproject.spi.transfer.provider.SignalHandler; +import org.datatransferproject.spi.transfer.provider.SignalRequest; +import org.datatransferproject.types.transfer.auth.AuthData; +import org.datatransferproject.types.transfer.auth.TokensAndUrlAuthData; +import org.datatransferproject.types.transfer.retry.RetryException; +import org.datatransferproject.types.transfer.retry.RetryStrategyLibrary; +import org.datatransferproject.types.transfer.retry.RetryingCallable; + +/** A {@link SignalHandler} for Synology. */ +public class SynologySignalHandler implements SignalHandler { + private final SynologyOAuthTokenManager tokenManager; + private final SynologyDTPService synologyDTPService; + private final RetryStrategyLibrary retryStrategyLibrary; + + /** + * Constructs a new {@code SynologySignalHandler} instance. + * + * @param synologyDTPService the Synology DTP service + * @param tokenManager the Synology OAuth token manager + * @param retryStrategyLibrary the retry strategy library for handling transient failures + */ + public SynologySignalHandler( + SynologyOAuthTokenManager tokenManager, + SynologyDTPService synologyDTPService, + RetryStrategyLibrary retryStrategyLibrary) { + this.tokenManager = tokenManager; + this.synologyDTPService = synologyDTPService; + this.retryStrategyLibrary = retryStrategyLibrary; + } + + @Override + public void sendSignal(SignalRequest signalRequest, AuthData authData, Monitor monitor) + throws RetryException { + Objects.requireNonNull(signalRequest, "signalRequest cannot be null"); + Objects.requireNonNull(authData, "authData cannot be null"); + Objects.requireNonNull(monitor, "monitor cannot be null"); + + UUID uuidJobId = UUID.fromString(signalRequest.jobId()); + tokenManager.addAuthDataIfNotExist(uuidJobId, (TokensAndUrlAuthData) authData); + + monitor.info( + () -> + String.format( + "[SynologySignalHandler] Received signal for jobId: %s, jobStatus.state: %s," + + " jobStatus.endReason: %s", + signalRequest.jobId(), + signalRequest.jobStatus().state(), + signalRequest.jobStatus().endReason())); + + Callable sendJobSignalCallable = + () -> { + Boolean success = + (Boolean) + this.synologyDTPService + .sendJobSignal(signalRequest.jobStatus(), uuidJobId) + .get("success"); + + if (success) { + monitor.debug( + () -> + String.format( + "[SynologySignalHandler] Successfully sent signal for jobId: %s," + + " jobStatus.state: %s, jobStatus.endReason: %s", + signalRequest.jobId(), + signalRequest.jobStatus().state(), + signalRequest.jobStatus().endReason())); + return null; + } + + throw new RuntimeException( + String.format( + "Failed to send signal for jobId: %s. Response with success=false.", + signalRequest.jobId())); + }; + + RetryingCallable retryingSendJobSignalCallable = + new RetryingCallable<>( + sendJobSignalCallable, retryStrategyLibrary, Clock.systemUTC(), monitor); + + try { + retryingSendJobSignalCallable.call(); + } catch (Throwable e) { + throw e; + } + } +} diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml b/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml index 07b57234b..aefd9d4f1 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml @@ -11,3 +11,4 @@ serviceConfig: createAlbum: "/import/album" uploadItem: "/import/item" addItemToAlbum: "/import/album/item" + signalJob: "/import/job/signal" diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java new file mode 100644 index 000000000..12f31498c --- /dev/null +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java @@ -0,0 +1,127 @@ +/* + * Copyright 2026 The Data Transfer Project Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.datatransferproject.datatransfer.synology.signals; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.Collections; +import java.util.Map; +import java.util.UUID; +import org.datatransferproject.api.launcher.Monitor; +import org.datatransferproject.datatransfer.synology.service.SynologyDTPService; +import org.datatransferproject.datatransfer.synology.service.SynologyOAuthTokenManager; +import org.datatransferproject.spi.transfer.provider.SignalRequest; +import org.datatransferproject.spi.transfer.types.CopyExceptionWithFailureReason; +import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle; +import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle.EndReason; +import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle.State; +import org.datatransferproject.types.common.models.DataVertical; +import org.datatransferproject.types.transfer.auth.TokensAndUrlAuthData; +import org.datatransferproject.types.transfer.retry.RetryException; +import org.datatransferproject.types.transfer.retry.RetryMapping; +import org.datatransferproject.types.transfer.retry.RetryStrategy; +import org.datatransferproject.types.transfer.retry.RetryStrategyLibrary; +import org.datatransferproject.types.transfer.retry.UniformRetryStrategy; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +public class SynologySignalHandlerTest { + + private SynologyOAuthTokenManager tokenManager; + private RetryStrategyLibrary retryStrategyLibrary; + private Monitor monitor; + private String jobId; + private SynologySignalHandler signalHandler; + private SynologyDTPService synologyDTPService; + private TokensAndUrlAuthData authData; + + @BeforeEach + public void setUp() { + jobId = UUID.randomUUID().toString(); + + tokenManager = mock(SynologyOAuthTokenManager.class); + RetryStrategy retryStrategy = new UniformRetryStrategy(3, 1, "test"); + RetryMapping retryMapping = + new RetryMapping( + new String[] {RuntimeException.class.getName()}, new String[] {}, retryStrategy); + retryStrategyLibrary = + new RetryStrategyLibrary( + Collections.singletonList(retryMapping), new UniformRetryStrategy(1, 1, "default")); + monitor = mock(Monitor.class); + synologyDTPService = mock(SynologyDTPService.class); + authData = mock(TokensAndUrlAuthData.class); + + signalHandler = + new SynologySignalHandler(tokenManager, synologyDTPService, retryStrategyLibrary); + } + + @Test + public void testSendSignal() throws RetryException, CopyExceptionWithFailureReason { + JobLifeCycle jobStatus = + JobLifeCycle.builder() + .setState(State.ENDED) + .setEndReason(EndReason.SUCCESSFULLY_COMPLETED) + .build(); + + SignalRequest signalRequest = + SignalRequest.builder() + .setJobId(jobId) + .setJobStatus(jobStatus) + .setExportingService("EXPORT_SERVICE") + .setImportingService("IMPORT_SERVICE") + .setDataType(DataVertical.MAIL.getDataType()) + .build(); + + when(synologyDTPService.sendJobSignal(any(JobLifeCycle.class), any(UUID.class))) + .thenReturn(Map.of("success", true)); + + signalHandler.sendSignal(signalRequest, authData, monitor); + + verify(synologyDTPService, times(1)).sendJobSignal(any(JobLifeCycle.class), any(UUID.class)); + } + + @Test + public void testSendSignalRetry() throws RetryException, CopyExceptionWithFailureReason { + JobLifeCycle jobStatus = + JobLifeCycle.builder() + .setState(State.ENDED) + .setEndReason(EndReason.SUCCESSFULLY_COMPLETED) + .build(); + + SignalRequest signalRequest = + SignalRequest.builder() + .setJobId(jobId) + .setJobStatus(jobStatus) + .setExportingService("EXPORT_SERVICE") + .setImportingService("IMPORT_SERVICE") + .setDataType(DataVertical.MAIL.getDataType()) + .build(); + + when(synologyDTPService.sendJobSignal(any(JobLifeCycle.class), any(UUID.class))) + .thenThrow(new RuntimeException("Failed to send signal")) + .thenReturn(Map.of("success", true)); + + signalHandler.sendSignal(signalRequest, authData, monitor); + + verify(synologyDTPService, times(2)).sendJobSignal(any(JobLifeCycle.class), any(UUID.class)); + } +} diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java index 87ea1bf05..fa2280a49 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java @@ -30,12 +30,16 @@ public class TestConfigs { public static final String TEST_CREATE_ALBUM_PATH = "/create"; public static final String TEST_UPLOAD_ITEM_PATH = "/upload"; public static final String TEST_ADD_ITEM_TO_ALBUM_PATH = "/add"; + public static final String TEST_SIGNAL_JOB_PATH = "/signal"; public static final int TEST_MAX_ATTEMPTS = 5; public static ServiceConfig createServiceConfig() { C2Api.ApiPath apiPath = new C2Api.ApiPath( - TEST_CREATE_ALBUM_PATH, TEST_UPLOAD_ITEM_PATH, TEST_ADD_ITEM_TO_ALBUM_PATH); + TEST_CREATE_ALBUM_PATH, + TEST_UPLOAD_ITEM_PATH, + TEST_ADD_ITEM_TO_ALBUM_PATH, + TEST_SIGNAL_JOB_PATH); C2Api c2Api = new C2Api(TEST_C2_BASE_URL, apiPath); From e0017d04e09e17443584248ac15a35042efcd2c5 Mon Sep 17 00:00:00 2001 From: "Chih-Hsien (Simon) Yeh" Date: Wed, 18 Mar 2026 16:04:06 +0800 Subject: [PATCH 03/11] feat(synology): implement streaming for large file uploads (#1485) ## Goal The goal of this change is to support large file uploads for `SynologyDTPService` by implementing streaming. This resolves potential OutOfMemoryError (OOM) issues that occurred when large media files were fully buffered into memory during the transfer process. ## Changes - **Streaming Uploads:** Replaced `ByteStreams.toByteArray()` with a custom `RequestBody` implementation using **Okio** to stream data directly from the source (`JobStore` or `URL`) to the network sink. - **Repeatable Streams for Retries:** Introduced `RequestBodyGenerator`, a functional interface that allows the `sendPostRequest` method to re-open the `InputStream` during retries. This ensures that even if a stream is consumed during a failed attempt, it can be reset for the next retry. ## Testing - **New Unit Tests:** Added `SynologyDTPServiceOOMTest` to verify 1GB streaming. - **Updated Unit Tests:** Updated `SynologyDTPServiceTest` to accommodate the new `RequestBodyGenerator` pattern and added cases for `getMediaInputStreamWrapper`. --- .../synology/service/SynologyDTPService.java | 318 ++++++++++++------ .../synology/uploader/SynologyUploader.java | 106 +++--- .../service/SynologyDTPServiceOOMTest.java | 227 +++++++++++++ .../service/SynologyDTPServiceTest.java | 307 +++++++++-------- .../signals/SynologySignalHandlerTest.java | 6 +- .../uploader/SynologyUploaderTest.java | 55 +-- 6 files changed, 671 insertions(+), 348 deletions(-) create mode 100644 extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java index e946e2fda..d69202585 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java @@ -20,11 +20,9 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Strings; -import com.google.common.io.ByteStreams; import java.io.IOException; -import java.io.InputStream; -import java.net.MalformedURLException; import java.net.URL; +import java.net.URLConnection; import java.util.Date; import java.util.Map; import java.util.UUID; @@ -35,15 +33,19 @@ import okhttp3.Request; import okhttp3.RequestBody; import okhttp3.Response; +import okio.BufferedSink; +import okio.Okio; +import okio.Source; import org.datatransferproject.api.launcher.Monitor; import org.datatransferproject.datatransfer.synology.models.C2Api; import org.datatransferproject.datatransfer.synology.models.ServiceConfig; import org.datatransferproject.datatransfer.synology.utils.ServiceConfigParser; import org.datatransferproject.spi.cloud.storage.JobStore; +import org.datatransferproject.spi.cloud.storage.TemporaryPerJobDataStore.InputStreamWrapper; import org.datatransferproject.spi.transfer.types.CopyExceptionWithFailureReason; import org.datatransferproject.spi.transfer.types.DestinationMemoryFullException; +import org.datatransferproject.spi.transfer.types.InvalidTokenException; import org.datatransferproject.spi.transfer.types.NoNasInAccountException; -import org.datatransferproject.spi.transfer.types.UploadErrorException; import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle; import org.datatransferproject.types.common.models.media.MediaAlbum; import org.datatransferproject.types.common.models.photos.PhotoModel; @@ -61,6 +63,11 @@ public class SynologyDTPService { C2Api c2Api; ServiceConfig.Retry retryConfig; + @FunctionalInterface + protected interface RequestBodyGenerator { + RequestBody get() throws CopyExceptionWithFailureReason, IOException; + } + /** * Constructs a new {@code SynologyDTPService} instance. * @@ -89,17 +96,6 @@ public SynologyDTPService( this.client = client; } - /** - * getInputStream - * - * @param url - * @return an InputStream Instance - * @throws IOException - */ - protected InputStream getInputStream(String url) throws IOException { - return new URL(url).openStream(); - } - /** * Creates album. * @@ -108,15 +104,38 @@ protected InputStream getInputStream(String url) throws IOException { * @return a map of shape {"data": {"album_id": }} */ public Map createAlbum(MediaAlbum album, UUID jobId) - throws CopyExceptionWithFailureReason { + throws CopyExceptionWithFailureReason, IOException { FormBody.Builder builder = new FormBody.Builder().add("title", album.getName()); builder.add("album_id", album.getId()); builder.add("job_id", jobId.toString()); builder.add("service", exportingService); monitor.info(() -> "[SynologyImporter] Creating album", album.getName(), jobId); + RequestBody requestBody = builder.build(); return (Map) - sendPostRequest(c2Api.getCreateAlbum(), builder.build(), jobId).get("data"); + sendPostRequest(c2Api.getCreateAlbum(), () -> requestBody, jobId).get("data"); + } + + /** + * get InputStreamWrapper for media file, it can be from temp store or from fetchable url + * + * @param jobId the job ID + * @param fetchableUrl the url to fetch media file, can be null if the file is in temp store + * @return an InputStreamWrapper instance + * @throws CopyExceptionWithFailureReason + */ + @VisibleForTesting + protected InputStreamWrapper getMediaInputStreamWrapper( + UUID jobId, String fetchableUrl, boolean isInTempStore) throws IOException { + if (isInTempStore) { + return jobStore.getStream(jobId, fetchableUrl); + } else if (fetchableUrl != null) { + URL url = new URL(fetchableUrl); + URLConnection connection = url.openConnection(); + return new InputStreamWrapper(connection.getInputStream(), connection.getContentLengthLong()); + } + + throw new IllegalArgumentException("fetchableUrl is null and isInTempStore is false"); } /** @@ -127,50 +146,78 @@ public Map createAlbum(MediaAlbum album, UUID jobId) * @return a map of shape {"data": {"item_id": }} */ public Map createPhoto(PhotoModel photo, UUID jobId) - throws CopyExceptionWithFailureReason { - byte[] imageBytes; - try { - InputStream inputStream = null; - if (photo.isInTempStore()) { - inputStream = jobStore.getStream(jobId, photo.getFetchableUrl()).getStream(); - } else if (photo.getFetchableUrl() != null) { - inputStream = getInputStream(photo.getFetchableUrl()); - } else { - monitor.severe(() -> "[SynologyImporter] Can't get inputStream for a photo"); - return null; - } - imageBytes = ByteStreams.toByteArray(inputStream); - } catch (MalformedURLException e) { - throw new UploadErrorException("Failed to create url for photo", e); - } catch (IOException e) { - throw new UploadErrorException("Failed to create input stream for photo", e); - } - - RequestBody fileBody = RequestBody.create(MediaType.parse(photo.getMimeType()), imageBytes); - - MultipartBody.Builder builder = - new MultipartBody.Builder() - .setType(MultipartBody.FORM) - .addFormDataPart("file", photo.getTitle(), fileBody) - .addFormDataPart("item_id", photo.getDataId()) - .addFormDataPart("title", photo.getTitle()) - .addFormDataPart("job_id", jobId.toString()) - .addFormDataPart("service", exportingService); + throws CopyExceptionWithFailureReason, IOException { + monitor.info( + () -> + String.format( + "[SynologyImporter] starts creating photo, dataId: [%s], name: [%s].", + photo.getDataId(), photo.getTitle()), + jobId); + + RequestBodyGenerator bodyGenerator = + () -> { + // Due to InputStream may not repeatable, we need to open it inside the generator function + // to make sure it can be read when retrying. + InputStreamWrapper inputStreamWrapper = + getMediaInputStreamWrapper(jobId, photo.getFetchableUrl(), photo.isInTempStore()); + + RequestBody fileBody = + new RequestBody() { + private boolean isConsumed = false; + + @Override + public MediaType contentType() { + return MediaType.parse(photo.getMimeType()); + } + + @Override + public long contentLength() { + return inputStreamWrapper.getBytes(); + } + + @Override + public void writeTo(BufferedSink sink) throws IOException { + if (isConsumed) { + throw new IOException("InputStream has already been consumed"); + } + isConsumed = true; + try (Source source = Okio.source(inputStreamWrapper.getStream())) { + sink.writeAll(source); + } + } + }; + + MultipartBody.Builder builder = + new MultipartBody.Builder() + .setType(MultipartBody.FORM) + .addFormDataPart("file", photo.getTitle(), fileBody) + .addFormDataPart("item_id", photo.getDataId()) + .addFormDataPart("title", photo.getTitle()) + .addFormDataPart("job_id", jobId.toString()) + .addFormDataPart("service", exportingService); + + String imageDescription = photo.getDescription(); + if (!Strings.isNullOrEmpty(imageDescription)) { + builder.addFormDataPart("description", imageDescription); + } + Date imageUploadedTime = photo.getUploadedTime(); + if (imageUploadedTime != null) { + long timestampInSeconds = imageUploadedTime.getTime() / 1000; + builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); + } - String imageDescription = photo.getDescription(); - if (!Strings.isNullOrEmpty(imageDescription)) { - builder.addFormDataPart("description", imageDescription); - } - Date imageUploadedTime = photo.getUploadedTime(); - if (imageUploadedTime != null) { - long timestampInSeconds = imageUploadedTime.getTime() / 1000; - builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); - } + return builder.build(); + }; @SuppressWarnings("unchecked") Map responseData = (Map) - sendPostRequest(c2Api.getUploadItem(), builder.build(), jobId).get("data"); + sendPostRequest(c2Api.getUploadItem(), bodyGenerator, jobId).get("data"); + monitor.info( + () -> + String.format( + "[SynologyImporter] photo created successfully, name: [%s].", photo.getTitle()), + jobId); return responseData; } @@ -182,50 +229,78 @@ public Map createPhoto(PhotoModel photo, UUID jobId) * @return a map of shape {"data": {"item_id": }} */ public Map createVideo(VideoModel video, UUID jobId) - throws CopyExceptionWithFailureReason { - byte[] videoBytes; - try { - InputStream inputStream = null; - if (video.isInTempStore()) { - inputStream = jobStore.getStream(jobId, video.getFetchableUrl()).getStream(); - } else if (video.getFetchableUrl() != null) { - inputStream = getInputStream(video.getFetchableUrl()); - } else { - monitor.severe(() -> "[SynologyImporter] Can't get inputStream for a video"); - return null; - } - - videoBytes = ByteStreams.toByteArray(inputStream); - } catch (MalformedURLException e) { - throw new UploadErrorException("Failed to create url for video", e); - } catch (IOException e) { - throw new UploadErrorException("Failed to create input stream for video", e); - } + throws CopyExceptionWithFailureReason, IOException { + monitor.info( + () -> + String.format( + "[SynologyImporter] starts creating video, dataId: [%s], name: [%s].", + video.getDataId(), video.getName()), + jobId); + + RequestBodyGenerator bodyGenerator = + () -> { + // Due to InputStream may not repeatable, we need to open it inside the generator function + // to make sure it can be read when retrying. + InputStreamWrapper inputStreamWrapper = + getMediaInputStreamWrapper(jobId, video.getFetchableUrl(), video.isInTempStore()); + + RequestBody fileBody = + new RequestBody() { + private boolean isConsumed = false; + + @Override + public MediaType contentType() { + return MediaType.parse(video.getMimeType()); + } + + @Override + public long contentLength() { + return inputStreamWrapper.getBytes(); + } + + @Override + public void writeTo(BufferedSink sink) throws IOException { + if (isConsumed) { + throw new IOException("InputStream has already been consumed"); + } + isConsumed = true; + try (Source source = Okio.source(inputStreamWrapper.getStream())) { + sink.writeAll(source); + } + } + }; + + MultipartBody.Builder builder = + new MultipartBody.Builder() + .setType(MultipartBody.FORM) + .addFormDataPart("file", video.getName(), fileBody) + .addFormDataPart("item_id", video.getDataId()) + .addFormDataPart("title", video.getName()) + .addFormDataPart("job_id", jobId.toString()) + .addFormDataPart("service", exportingService); + + String imageDescription = video.getDescription(); + if (!Strings.isNullOrEmpty(imageDescription)) { + builder.addFormDataPart("description", imageDescription); + } + Date videoUploadedTime = video.getUploadedTime(); + if (videoUploadedTime != null) { + long timestampInSeconds = videoUploadedTime.getTime() / 1000; + builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); + } - RequestBody fileBody = RequestBody.create(MediaType.parse(video.getMimeType()), videoBytes); - MultipartBody.Builder builder = - new MultipartBody.Builder() - .setType(MultipartBody.FORM) - .addFormDataPart("file", video.getName(), fileBody) - .addFormDataPart("item_id", video.getDataId()) - .addFormDataPart("title", video.getName()) - .addFormDataPart("job_id", jobId.toString()) - .addFormDataPart("service", exportingService); - - String imageDescription = video.getDescription(); - if (!Strings.isNullOrEmpty(imageDescription)) { - builder.addFormDataPart("description", imageDescription); - } - Date videoUploadedTime = video.getUploadedTime(); - if (videoUploadedTime != null) { - long timestampInSeconds = videoUploadedTime.getTime() / 1000; - builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); - } + return builder.build(); + }; @SuppressWarnings("unchecked") Map responseData = (Map) - sendPostRequest(c2Api.getUploadItem(), builder.build(), jobId).get("data"); + sendPostRequest(c2Api.getUploadItem(), bodyGenerator, jobId, 300).get("data"); + monitor.info( + () -> + String.format( + "[SynologyImporter] video created successfully, name: [%s].", video.getName()), + jobId); return responseData; } @@ -237,14 +312,15 @@ public Map createVideo(VideoModel video, UUID jobId) * @return a map of shape {"success": } */ public Map addItemToAlbum(String albumId, String itemId, UUID jobId) - throws CopyExceptionWithFailureReason { + throws CopyExceptionWithFailureReason, IOException { FormBody.Builder builder = new FormBody.Builder() .add("job_id", jobId.toString()) .add("service", exportingService) .add("album_id", albumId) .add("item_id", itemId); - return sendPostRequest(c2Api.getAddItemToAlbum(), builder.build(), jobId); + RequestBody requestBody = builder.build(); + return sendPostRequest(c2Api.getAddItemToAlbum(), () -> requestBody, jobId); } /** @@ -255,7 +331,7 @@ public Map addItemToAlbum(String albumId, String itemId, UUID jo * @return a map of shape {"success": } */ public Map sendJobSignal(JobLifeCycle jobStatus, UUID jobId) - throws CopyExceptionWithFailureReason { + throws CopyExceptionWithFailureReason, IOException { FormBody.Builder builder = new FormBody.Builder() .add("job_id", jobId.toString()) @@ -265,7 +341,8 @@ public Map sendJobSignal(JobLifeCycle jobStatus, UUID jobId) builder.add("end_reason", jobStatus.endReason().name()); } - return sendPostRequest(c2Api.getSignalJob(), builder.build(), jobId); + RequestBody requestBody = builder.build(); + return sendPostRequest(c2Api.getSignalJob(), () -> requestBody, jobId); } @VisibleForTesting @@ -304,20 +381,51 @@ protected void throwExceptionIfNoQuota(Response response) throws CopyExceptionWi } @VisibleForTesting - protected Map sendPostRequest(String url, RequestBody body, UUID jobId) - throws CopyExceptionWithFailureReason { + protected Map sendPostRequest( + String url, RequestBodyGenerator bodyGenerator, UUID jobId) + throws CopyExceptionWithFailureReason, IOException { + return sendPostRequest(url, bodyGenerator, jobId, -1); + } + + /* + * @param url the URL to send the POST request to + * @param bodyGenerator a generator function that produces the request body, it will be called for each retry attempt + * @param jobId the job ID + * @param timeoutInSeconds the timeout for the request in seconds, -1 means do not modify the default timeout of OkHttpClient + */ + @VisibleForTesting + protected Map sendPostRequest( + String url, RequestBodyGenerator bodyGenerator, UUID jobId, int timeoutInSeconds) + throws CopyExceptionWithFailureReason, IOException { boolean triedRefreshToken = false; - Request.Builder requestBuilder = new Request.Builder().url(url).post(body); Exception lastException = null; for (int retry = retryConfig.getMaxAttempts(); retry > 0; retry--) { Response response = null; try { + RequestBody body = bodyGenerator.get(); + Request.Builder requestBuilder = new Request.Builder().url(url).post(body); requestBuilder.header("Authorization", "Bearer " + tokenManager.getAccessToken(jobId)); - response = client.newCall(requestBuilder.build()).execute(); + if (timeoutInSeconds < 0) { + // Use default timeout of OkHttpClient + response = client.newCall(requestBuilder.build()).execute(); + } else { + response = + client + .newBuilder() + .readTimeout(timeoutInSeconds, java.util.concurrent.TimeUnit.SECONDS) + .build() + .newCall(requestBuilder.build()) + .execute(); + } if (!response.isSuccessful()) { int code = response.code(); - if (code == 401 && !triedRefreshToken) { + if (code == 401) { + if (triedRefreshToken) { + throw new InvalidTokenException( + "Synology access token is invalid even after refresh", + new IOException("SynologyDTPService get http 401 unauthorized")); + } triedRefreshToken = true; if (tokenManager.refreshToken(jobId, client, objectMapper)) { continue; @@ -367,7 +475,7 @@ protected Map sendPostRequest(String url, RequestBody body, UUID lastException = e; } } - throw new UploadErrorException( + throw new IOException( String.format( "Failed to send POST request after %d retries: %s", retryConfig.getMaxAttempts(), lastException.getMessage()), diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploader.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploader.java index eb3b343f3..de0277c66 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploader.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploader.java @@ -17,6 +17,7 @@ package org.datatransferproject.datatransfer.synology.uploader; +import java.io.IOException; import java.util.Collection; import java.util.Map; import java.util.UUID; @@ -27,8 +28,6 @@ import org.datatransferproject.datatransfer.synology.utils.SynologyMediaAlbumBinder; import org.datatransferproject.spi.transfer.idempotentexecutor.IdempotentImportExecutor; import org.datatransferproject.spi.transfer.types.CopyExceptionWithFailureReason; -import org.datatransferproject.spi.transfer.types.UploadErrorException; -import org.datatransferproject.types.common.DownloadableItem; import org.datatransferproject.types.common.ImportableItem; import org.datatransferproject.types.common.models.media.MediaAlbum; import org.datatransferproject.types.common.models.photos.PhotoAlbum; @@ -72,10 +71,13 @@ public SynologyUploader( */ public void importAlbums(Collection albums, UUID jobId) throws CopyExceptionWithFailureReason { + monitor.info( + () -> String.format("[SynologyImporter] starts importing albums, size: %d.", albums.size()), + jobId); if (albums.isEmpty()) { return; } - monitor.info(() -> "[SynologyImporter] starts importing albums", jobId); + for (ImportableItem album : albums) { try { MediaAlbum mediaAlbum; @@ -86,9 +88,13 @@ public void importAlbums(Collection albums, UUID jobId } else if (album instanceof VideoAlbum) { mediaAlbum = MediaAlbum.videoToMediaAlbum(((VideoAlbum) album)); } else { - throw new UploadErrorException( - "cannot convert to MediaAlbum", - new IllegalArgumentException("Unsupported ImportableItem type: " + album.getClass())); + monitor.severe( + () -> + String.format( + "[SynologyImporter] unsupported album type: %s, album name: %s.", + album.getClass(), album.getName()), + jobId); + continue; } String newAlbumId = importItemWithCache( @@ -104,10 +110,9 @@ public void importAlbums(Collection albums, UUID jobId throw e; } catch (Exception e) { monitor.severe(e::toString, jobId); - throw new UploadErrorException("Failed to import albums", e); } } - monitor.info(() -> "[SynologyImporter] imported albums successfully", jobId); + monitor.info(() -> "[SynologyImporter] ended importing albums.", jobId); } /** @@ -118,17 +123,26 @@ public void importAlbums(Collection albums, UUID jobId */ public void importPhotos(Collection photos, UUID jobId) throws CopyExceptionWithFailureReason { + monitor.info( + () -> String.format("[SynologyImporter] starts importing photos, size: %d.", photos.size()), + jobId); if (photos.isEmpty()) { return; } - monitor.info(() -> "[SynologyImporter] starts importing photos", jobId); + monitor.info(() -> "[SynologyImporter] starts importing photos.", jobId); for (PhotoModel photo : photos) { + monitor.info( + () -> + String.format( + "[SynologyImporter] starts importing photo, dataId: [%s], name: [%s]", + photo.getDataId(), photo.getName()), + jobId); try { String newPhotoId = - importDownloadableItemWithCache( + importItemWithCache( photo, jobId, - PhotoModel::getAlbumId, + "item_id", synologyDTPService::createPhoto, PhotoModel::getDataId, PhotoModel::getName); @@ -138,10 +152,9 @@ public void importPhotos(Collection photos, UUID jobId) throw e; } catch (Exception e) { monitor.severe(e::toString, jobId); - throw new UploadErrorException("Failed to import photos", e); } } - monitor.info(() -> "[SynologyImporter] imported photos successfully", jobId); + monitor.info(() -> "[SynologyImporter] ended importing photos.", jobId); } /** @@ -152,17 +165,26 @@ public void importPhotos(Collection photos, UUID jobId) */ public void importVideos(Collection videos, UUID jobId) throws CopyExceptionWithFailureReason { + monitor.info( + () -> String.format("[SynologyImporter] starts importing videos, size: %d", videos.size()), + jobId); if (videos.isEmpty()) { return; } monitor.info(() -> "[SynologyImporter] starts importing videos", jobId); for (VideoModel video : videos) { + monitor.info( + () -> + String.format( + "[SynologyImporter] starts importing video, dataId: [%s], name: [%s]", + video.getDataId(), video.getName()), + jobId); try { String newVideoId = - importDownloadableItemWithCache( + importItemWithCache( video, jobId, - VideoModel::getAlbumId, + "item_id", synologyDTPService::createVideo, VideoModel::getDataId, VideoModel::getName); @@ -172,46 +194,14 @@ public void importVideos(Collection videos, UUID jobId) throw e; } catch (Exception e) { monitor.severe(e::toString, jobId); - throw new UploadErrorException("Failed to import videos", e); } } - monitor.info(() -> "[SynologyImporter] imported videos successfully", jobId); + monitor.info(() -> "[SynologyImporter] ended importing videos.", jobId); } @FunctionalInterface private interface SynologyBiFunction { - R apply(T t, U u) throws CopyExceptionWithFailureReason; - } - - /** - * Imports a item - * - * @param item the item - * @return the item ID - */ - private String importDownloadableItemWithCache( - T item, - UUID jobId, - Function albumIdFunction, - SynologyBiFunction> createFunction, - Function dataIdFunction, - Function nameFunction) - throws CopyExceptionWithFailureReason { - String newItemId = null; - try { - newItemId = - importItemWithCache(item, jobId, "item_id", createFunction, dataIdFunction, nameFunction); - } catch (CopyExceptionWithFailureReason e) { - throw e; - } catch (Exception e) { - throw new UploadErrorException( - String.format( - "Failed to import item [%s], name [%s]", - dataIdFunction.apply(item), nameFunction.apply(item)), - e); - } - - return newItemId; + R apply(T t, U u) throws CopyExceptionWithFailureReason, IOException; } private String importItemWithCache( @@ -250,16 +240,22 @@ private void addItemToAlbum(String albumId, String itemId, UUID jobId) (Boolean) synologyDTPService.addItemToAlbum(albumId, itemId, jobId).get("success")); if (Boolean.FALSE.equals(createResult)) { - throw new UploadErrorException( - String.format( - "Unsuccessful result from adding item [%s] to album [%s]", itemId, albumId), - null); + monitor.severe( + () -> + String.format( + "[SynologyImporter] failed to add item to album, albumId: [%s], itemId: [%s].", + albumId, itemId)); } } catch (CopyExceptionWithFailureReason e) { throw e; } catch (Exception e) { - throw new UploadErrorException( - String.format("Failed to add item [%s] to album [%s]", itemId, albumId), e); + monitor.severe( + () -> + String.format( + "[SynologyImporter] failed to add item to album, albumId: [%s], itemId: [%s].", + albumId, itemId), + e); + return; } monitor.info( () -> "[SynologyImporter] added item to album successfully", diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java new file mode 100644 index 000000000..48b5510cd --- /dev/null +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java @@ -0,0 +1,227 @@ +/* + * Copyright 2025 The Data Transfer Project Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.datatransferproject.datatransfer.synology.service; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.when; + +import java.io.IOException; +import java.io.InputStream; +import java.util.Map; +import java.util.UUID; +import okhttp3.OkHttpClient; +import okhttp3.RequestBody; +import okio.Buffer; +import okio.Okio; +import okio.Sink; +import okio.Timeout; +import org.datatransferproject.api.launcher.Monitor; +import org.datatransferproject.datatransfer.synology.utils.TestConfigs; +import org.datatransferproject.spi.cloud.storage.JobStore; +import org.datatransferproject.spi.cloud.storage.TemporaryPerJobDataStore.InputStreamWrapper; +import org.datatransferproject.spi.transfer.types.InvalidTokenException; +import org.datatransferproject.types.common.models.photos.PhotoModel; +import org.datatransferproject.types.common.models.videos.VideoModel; +import org.datatransferproject.types.transfer.serviceconfig.TransferServiceConfig; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +public class SynologyDTPServiceOOMTest { + private final String exportingService = "mockService"; + private final UUID jobId = UUID.randomUUID(); + private SynologyDTPService dtpService; + @Mock private Monitor monitor; + @Mock private TransferServiceConfig transferServiceConfig; + @Mock private JobStore jobStore; + @Mock private SynologyOAuthTokenManager tokenManager; + @Captor private ArgumentCaptor requestBodyCaptor; + @Mock private OkHttpClient client; + + // Helper class for OOM test + private static class FakeLargeInputStream extends InputStream { + private final long size; + private long bytesRead = 0; + private final byte[] singleByte = new byte[] {'a'}; + + public FakeLargeInputStream(long size) { + this.size = size; + } + + @Override + public int read() { + if (bytesRead >= size) { + return -1; // End of stream + } + bytesRead++; + return singleByte[0]; + } + + @Override + public int read(byte[] b, int off, int len) { + if (b == null) { + throw new NullPointerException(); + } else if (off < 0 || len < 0 || len > b.length - off) { + throw new IndexOutOfBoundsException(); + } else if (len == 0) { + return 0; + } + + if (bytesRead >= size) { + return -1; + } + long remaining = size - bytesRead; + int toRead = (int) Math.min(len, remaining); + + // Don't bother filling the array, we are just simulating reading + // Arrays.fill(b, off, off + toRead, singleByte[0]); + + bytesRead += toRead; + return toRead; + } + + @Override + public int available() { + long remaining = size - bytesRead; + return remaining > Integer.MAX_VALUE ? Integer.MAX_VALUE : (int) remaining; + } + } + + private static class BlackholeSink implements Sink { + @Override + public void write(Buffer source, long byteCount) throws IOException { + source.skip(byteCount); + } + + @Override + public void flush() throws IOException {} + + @Override + public Timeout timeout() { + return Timeout.NONE; + } + + @Override + public void close() throws IOException {} + } + + @BeforeEach + public void setUp() throws InvalidTokenException { + when(transferServiceConfig.getServiceConfig()) + .thenReturn(TestConfigs.createServiceConfigJson()); + dtpService = + new SynologyDTPService( + monitor, transferServiceConfig, exportingService, jobStore, tokenManager, client); + } + + @Test + public void createPhoto_withLargeFile_shouldStreamData() throws Exception { + // setup + long oneGB = 1024L * 1024L * 1024L; + InputStream fakeInputStream = new FakeLargeInputStream(oneGB); + InputStreamWrapper streamWrapper = new InputStreamWrapper(fakeInputStream, oneGB); + + PhotoModel photo = + new PhotoModel( + "large-photo", "large-photo-url", "desc", "image/jpeg", "photo-id", null, true); + + when(jobStore.getStream(jobId, "large-photo-url")).thenReturn(streamWrapper); + + SynologyDTPService spyService = Mockito.spy(dtpService); + doReturn(Map.of("success", true, "data", Map.of("item_id", "photo-id"))) + .when(spyService) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + + // act + spyService.createPhoto(photo, jobId); + + // Get the generated request body + RequestBody requestBody = requestBodyCaptor.getValue().get(); + + // assert + // The content length is not asserted to be equal to 1GB because it includes multipart + // boundaries and headers, + // making it larger than the raw file size. The main purpose of this test is to ensure that + // the file is streamed without causing an OutOfMemoryError. + assertTrue(requestBody.contentLength() > oneGB); + + // Simulate writing the body to a sink that discards the data. + // This will throw OutOfMemoryError if the whole stream is loaded into memory. + okio.BufferedSink discardingSink = Okio.buffer(new BlackholeSink()); + + // The MultipartBody will try to read the stream to write it. + // If it buffers the whole 1GB in memory, this will OOM. + requestBody.writeTo(discardingSink); + discardingSink.flush(); + + // If we reach here, it means we streamed the data without OOM. + } + + @Test + public void createVideo_withLargeFile_shouldStreamData() throws Exception { + // setup + long oneGB = 1024L * 1024L * 1024L; + InputStream fakeInputStream = new FakeLargeInputStream(oneGB); + InputStreamWrapper streamWrapper = new InputStreamWrapper(fakeInputStream, oneGB); + + VideoModel video = + new VideoModel( + "large-video", "large-video-url", "desc", "video/mp4", "video-id", null, true, null); + + when(jobStore.getStream(jobId, "large-video-url")).thenReturn(streamWrapper); + + SynologyDTPService spyService = Mockito.spy(dtpService); + doReturn(Map.of("success", true, "data", Map.of("item_id", "video-id"))) + .when(spyService) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); + + // act + spyService.createVideo(video, jobId); + + // Get the generated request body + RequestBody requestBody = requestBodyCaptor.getValue().get(); + + // assert + // The content length is not asserted to be equal to 1GB because it includes multipart + // boundaries and headers, + // making it larger than the raw file size. The main purpose of this test is to ensure that + // the file is streamed without causing an OutOfMemoryError. + assertTrue(requestBody.contentLength() > oneGB); + + // Simulate writing the body to a sink that discards the data. + // This will throw OutOfMemoryError if the whole stream is loaded into memory. + okio.BufferedSink discardingSink = Okio.buffer(new BlackholeSink()); + + // The MultipartBody will try to read the stream to write it. + // If it buffers the whole 1GB in memory, this will OOM. + requestBody.writeTo(discardingSink); + discardingSink.flush(); + + // If we reach here, it means we streamed the data without OOM. + } +} diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java index ba2925a65..cb6757905 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java @@ -24,14 +24,26 @@ import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; -import static org.mockito.Mockito.*; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.anyString; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.eq; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.InputStream; -import java.net.MalformedURLException; +import java.nio.file.Files; +import java.nio.file.Path; import java.util.Date; import java.util.HashMap; import java.util.Map; @@ -48,6 +60,7 @@ import okhttp3.ResponseBody; import okio.Buffer; import org.datatransferproject.api.launcher.Monitor; +import org.datatransferproject.datatransfer.synology.service.SynologyDTPService.RequestBodyGenerator; import org.datatransferproject.datatransfer.synology.utils.TestConfigs; import org.datatransferproject.spi.cloud.storage.JobStore; import org.datatransferproject.spi.cloud.storage.TemporaryPerJobDataStore.InputStreamWrapper; @@ -84,7 +97,7 @@ public class SynologyDTPServiceTest { @Mock protected TransferServiceConfig transferServiceConfig; @Mock protected JobStore jobStore; @Mock protected SynologyOAuthTokenManager tokenManager; - @Captor ArgumentCaptor requestBodyCaptor; + @Captor ArgumentCaptor requestBodyCaptor; @Mock private OkHttpClient client; @BeforeEach @@ -105,16 +118,19 @@ public class AddItemToAlbum { private final String itemId = "testItem"; @Test - public void shouldSendPostRequestWithCorrectFormBody() throws CopyExceptionWithFailureReason { + public void shouldSendPostRequestWithCorrectFormBody() + throws CopyExceptionWithFailureReason, IOException { SynologyDTPService spyService = Mockito.spy(dtpService); - doReturn(Map.of("success", true)).when(spyService).sendPostRequest(anyString(), any(), any()); + doReturn(Map.of("success", true)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); Map result = spyService.addItemToAlbum(albumId, itemId, jobId); verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); assertEquals(result.get("success"), true); - RequestBody capturedBody = requestBodyCaptor.getValue(); + RequestBody capturedBody = requestBodyCaptor.getValue().get(); assertTrue(capturedBody instanceof FormBody); FormBody formBody = (FormBody) capturedBody; @@ -131,12 +147,12 @@ public void shouldSendPostRequestWithCorrectFormBody() throws CopyExceptionWithF @Test public void shouldThrowExceptionIfSendPostRequestFailed() - throws CopyExceptionWithFailureReason { + throws CopyExceptionWithFailureReason, IOException { SynologyDTPService spyService = Mockito.spy(dtpService); doThrow(new UploadErrorException("MockException", null)) .when(spyService) - .sendPostRequest(anyString(), any(), any()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); assertThrows( UploadErrorException.class, @@ -152,20 +168,21 @@ public class CreateAlbum { private final MediaAlbum album = new MediaAlbum(albumId, albumName, ""); @Test - public void shouldSendPostRequestWithCorrectFormBody() throws CopyExceptionWithFailureReason { + public void shouldSendPostRequestWithCorrectFormBody() + throws CopyExceptionWithFailureReason, IOException { SynologyDTPService spyService = Mockito.spy(dtpService); Map dataMap = Map.of("album_id", albumId); doReturn(Map.of("success", true, "data", dataMap)) .when(spyService) - .sendPostRequest(anyString(), any(), any()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); Map result = spyService.createAlbum(album, jobId); verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); assertEquals(result.get("success"), null); assertEquals(result.get("album_id"), albumId); - RequestBody capturedBody = requestBodyCaptor.getValue(); + RequestBody capturedBody = requestBodyCaptor.getValue().get(); assertTrue(capturedBody instanceof FormBody); FormBody formBody = (FormBody) capturedBody; @@ -182,12 +199,12 @@ public void shouldSendPostRequestWithCorrectFormBody() throws CopyExceptionWithF @Test public void shouldThrowExceptionIfSendPostRequestFailed() - throws CopyExceptionWithFailureReason { + throws CopyExceptionWithFailureReason, IOException { SynologyDTPService spyService = Mockito.spy(dtpService); doThrow(new UploadErrorException("MockException", null)) .when(spyService) - .sendPostRequest(anyString(), any(), any()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); assertThrows( UploadErrorException.class, () -> spyService.createAlbum(album, jobId), "MockException"); @@ -224,108 +241,6 @@ public Stream provideMediaObjectsNotInTempStore() { new VideoModel(itemName, fetchUrl, description, "format", itemId, null, false, null)); } - @ParameterizedTest(name = "shouldGetStreamFromJobStoreIfMediaItemIsInTempStore [{index}] {0}") - @MethodSource("provideMediaObjectsInTempStore") - public void shouldGetStreamFromJobStoreIfMediaItemIsInTempStore(DownloadableFile item) - throws IOException, CopyExceptionWithFailureReason { - byte[] mockImage = new byte[] {1, 2, 3}; - InputStream mockInputStream = new ByteArrayInputStream(mockImage); - InputStreamWrapper streamWrapper = mock(InputStreamWrapper.class); - SynologyDTPService spyService = Mockito.spy(dtpService); - - when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); - when(streamWrapper.getStream()).thenReturn(mockInputStream); - doReturn(mock(Map.class)).when(spyService).sendPostRequest(anyString(), any(), any()); - - if (item instanceof PhotoModel) { - spyService.createPhoto((PhotoModel) item, jobId); - } else if (item instanceof VideoModel) { - spyService.createVideo((VideoModel) item, jobId); - } - - verify(spyService).sendPostRequest(anyString(), any(), any()); - verify(streamWrapper).getStream(); - verify(spyService, never()).getInputStream(fetchUrl); - } - - @ParameterizedTest(name = "shouldGetStreamFromUrlIfPhotoIsNotInTempStore [{index}] {0}") - @MethodSource("provideMediaObjectsNotInTempStore") - public void shouldGetStreamFromUrlIfPhotoIsNotInTempStore(DownloadableFile item) - throws IOException, CopyExceptionWithFailureReason { - byte[] mockImage = new byte[] {1, 2, 3}; - SynologyDTPService spyService = Mockito.spy(dtpService); - InputStream mockInputStream = new ByteArrayInputStream(mockImage); - - doReturn(mockInputStream).when(spyService).getInputStream(fetchUrl); - doReturn(mock(Map.class)).when(spyService).sendPostRequest(anyString(), any(), any()); - - if (item instanceof PhotoModel) { - spyService.createPhoto((PhotoModel) item, jobId); - } else if (item instanceof VideoModel) { - spyService.createVideo((VideoModel) item, jobId); - } - - verify(spyService).sendPostRequest(anyString(), any(), any()); - verify(spyService).getInputStream(fetchUrl); - verifyNoInteractions(jobStore); - } - - // VideoModel.contentUrl can't be null - @Test - public void shouldReturnNullWhenPhotoIsNotInTempStoreAndHasNoUrl() - throws CopyExceptionWithFailureReason { - PhotoModel photo = - new PhotoModel(itemName, null, description, "mediaType", itemId, null, false); - SynologyDTPService spyService = Mockito.spy(dtpService); - assertEquals(null, spyService.createPhoto((PhotoModel) photo, jobId)); - verify(spyService, never()).sendPostRequest(anyString(), any(), any()); - } - - @ParameterizedTest(name = "shouldThrowExceptionIfNewURLFailed [{index}] {0}") - @MethodSource("provideMediaObjectsNotInTempStore") - public void shouldThrowExceptionIfNewURLFailed(DownloadableFile item) - throws IOException, CopyExceptionWithFailureReason { - SynologyDTPService spyService = Mockito.spy(dtpService); - doThrow(new MalformedURLException("Failed to create url for photo")) - .when(spyService) - .getInputStream(fetchUrl); - - if (item instanceof PhotoModel) { - assertThrows( - UploadErrorException.class, - () -> spyService.createPhoto((PhotoModel) item, jobId), - "Failed to create url for photo"); - } else if (item instanceof VideoModel) { - assertThrows( - UploadErrorException.class, - () -> spyService.createVideo((VideoModel) item, jobId), - "Failed to create url for video"); - } - - verify(spyService, never()).sendPostRequest(anyString(), any(), any()); - } - - @ParameterizedTest(name = "shouldThrowExceptionIfFailedToGetStream [{index}] {0}") - @MethodSource("provideMediaObjectsInTempStore") - public void shouldThrowExceptionIfFailedToGetStream(DownloadableFile item) - throws IOException, CopyExceptionWithFailureReason { - SynologyDTPService spyService = Mockito.spy(dtpService); - when(jobStore.getStream(jobId, fetchUrl)).thenThrow(new IOException("Failed to get stream")); - - if (item instanceof PhotoModel) { - assertThrows( - UploadErrorException.class, - () -> spyService.createPhoto((PhotoModel) item, jobId), - "Failed to create input stream for photo"); - } else if (item instanceof VideoModel) { - assertThrows( - UploadErrorException.class, - () -> spyService.createVideo((VideoModel) item, jobId), - "Failed to create input stream for video"); - } - verify(spyService, never()).sendPostRequest(anyString(), any(), any()); - } - @ParameterizedTest( name = "shouldSendPostRequestWithCorrectFormBodyWithDescriptionAndUploadedTime [{index}] {0}") @@ -340,22 +255,26 @@ public void shouldSendPostRequestWithCorrectFormBodyWithDescriptionAndUploadedTi when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); when(streamWrapper.getStream()).thenReturn(mockInputStream); - doReturn(Map.of("success", true, "data", dataMap)) - .when(spyService) - .sendPostRequest(anyString(), any(), any()); - Map result = new HashMap(); if (item instanceof PhotoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createPhoto((PhotoModel) item, jobId); + verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); } else if (item instanceof VideoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); result = spyService.createVideo((VideoModel) item, jobId); + verify(spyService) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); } - verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); assertEquals(result.get("item_id"), itemId); assertEquals(result.get("success"), null); - RequestBody capturedBody = requestBodyCaptor.getValue(); + RequestBody capturedBody = requestBodyCaptor.getValue().get(); MultipartBody multipartBody = (MultipartBody) capturedBody; Map multipartFormAnswer = @@ -401,22 +320,26 @@ public void shouldSendPostRequestWithCorrectFormBodyWithoutDescription(Downloada when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); when(streamWrapper.getStream()).thenReturn(mockInputStream); - doReturn(Map.of("success", true, "data", dataMap)) - .when(spyService) - .sendPostRequest(anyString(), any(), any()); - Map result = new HashMap(); if (item instanceof PhotoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createPhoto((PhotoModel) item, jobId); + verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); } else if (item instanceof VideoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); result = spyService.createVideo((VideoModel) item, jobId); + verify(spyService) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); } - verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); assertEquals(result.get("item_id"), itemId); assertEquals(result.get("success"), null); - RequestBody capturedBody = requestBodyCaptor.getValue(); + RequestBody capturedBody = requestBodyCaptor.getValue().get(); MultipartBody multipartBody = (MultipartBody) capturedBody; Map multipartFormAnswer = @@ -446,27 +369,61 @@ public void shouldSendPostRequestWithCorrectFormBodyWithoutDescription(Downloada } } - @ParameterizedTest(name = "shouldThrowExceptionIfSendPostRequestFailed [{index}] {0}") + @ParameterizedTest(name = "shouldThrowExceptionIfInputStreamIsConsumed [{index}] {0}") @MethodSource("provideMediaObjectsInTempStore") - public void shouldThrowExceptionIfSendPostRequestFailed(DownloadableFile item) + public void shouldThrowExceptionIfInputStreamIsConsumed(DownloadableFile item) throws IOException, CopyExceptionWithFailureReason { byte[] mockImage = new byte[] {1, 2, 3}; InputStream mockInputStream = new ByteArrayInputStream(mockImage); InputStreamWrapper streamWrapper = mock(InputStreamWrapper.class); SynologyDTPService spyService = Mockito.spy(dtpService); + Map dataMap = Map.of("item_id", itemId); when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); when(streamWrapper.getStream()).thenReturn(mockInputStream); - doThrow(new UploadErrorException("MockException", null)) - .when(spyService) - .sendPostRequest(anyString(), any(), any()); + if (item instanceof PhotoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); + spyService.createPhoto((PhotoModel) item, jobId); + verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + } else if (item instanceof VideoModel) { + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); + spyService.createVideo((VideoModel) item, jobId); + verify(spyService) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); + } + + RequestBody capturedBody = requestBodyCaptor.getValue().get(); + Buffer buffer = new Buffer(); + capturedBody.writeTo(buffer); + IOException exception = assertThrows(IOException.class, () -> capturedBody.writeTo(buffer)); + assertTrue(exception.getMessage().contains("InputStream has already been consumed")); + } + + @ParameterizedTest(name = "shouldThrowExceptionIfSendPostRequestFailed [{index}] {0}") + @MethodSource("provideMediaObjectsInTempStore") + public void shouldThrowExceptionIfSendPostRequestFailed(DownloadableFile item) + throws IOException, CopyExceptionWithFailureReason { + byte[] mockImage = new byte[] {1, 2, 3}; + InputStream mockInputStream = new ByteArrayInputStream(mockImage); + InputStreamWrapper streamWrapper = mock(InputStreamWrapper.class); + SynologyDTPService spyService = Mockito.spy(dtpService); if (item instanceof PhotoModel) { + doThrow(new UploadErrorException("MockException", null)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); assertThrows( UploadErrorException.class, () -> spyService.createPhoto((PhotoModel) item, jobId), "MockException"); } else if (item instanceof VideoModel) { + doThrow(new UploadErrorException("MockException", null)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); assertThrows( UploadErrorException.class, () -> spyService.createVideo((VideoModel) item, jobId), @@ -497,7 +454,7 @@ public void shouldRetryIfGotError() throws IOException, CopyExceptionWithFailure when(mockCall.execute()).thenReturn(mockResponseFail).thenReturn(mockResponseSuccess); Map result = - dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId); + dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId); assertEquals(Map.of("success", true), result); @@ -518,9 +475,10 @@ public void shouldThrowExceptionIfReachMaxRetries() throws IOException { when(mockCall.execute()).thenReturn(mockResponseFail); assertThrows( - UploadErrorException.class, - () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId), - String.format("Failed to send POST request %d times", TestConfigs.TEST_MAX_ATTEMPTS)); + IOException.class, + () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId), + String.format( + "Failed to send POST request after %d retries", TestConfigs.TEST_MAX_ATTEMPTS)); verify(tokenManager, never()) .refreshToken(any(UUID.class), eq(client), any(ObjectMapper.class)); verify(client, times(TestConfigs.TEST_MAX_ATTEMPTS)).newCall(any(Request.class)); @@ -547,7 +505,7 @@ public void shouldRefreshTokenAndRetryIfGotUnauthorizedHttpError() .thenReturn(true); Map result = - dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId); + dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId); assertEquals(Map.of("success", true), result); @@ -571,13 +529,12 @@ public void shouldInvokeRefreshTokenOnlyOnceIfGotUnauthorizedMultipleTimes() .thenReturn(false); assertThrows( - UploadErrorException.class, - () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId), - String.format("Failed to send POST request %d times", TestConfigs.TEST_MAX_ATTEMPTS)); + InvalidTokenException.class, + () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId)); verify(tokenManager, times(1)) .refreshToken(any(UUID.class), eq(client), any(ObjectMapper.class)); - verify(client, times(TestConfigs.TEST_MAX_ATTEMPTS)).newCall(any(Request.class)); + verify(client, times(2)).newCall(any(Request.class)); } @Test @@ -594,9 +551,10 @@ public void shouldThrowExceptionIfParseResponseFailed() throws IOException { .thenThrow(new IOException("Error when call response.body.string()")); assertThrows( - UploadErrorException.class, - () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId), - String.format("Failed to send POST request %d times", TestConfigs.TEST_MAX_ATTEMPTS)); + IOException.class, + () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId), + String.format( + "Failed to send POST request after %d retries", TestConfigs.TEST_MAX_ATTEMPTS)); verify(tokenManager, never()) .refreshToken(any(UUID.class), eq(client), any(ObjectMapper.class)); verify(client, times(TestConfigs.TEST_MAX_ATTEMPTS)).newCall(any(Request.class)); @@ -616,7 +574,7 @@ public void shouldReturnResponseData() when(mockResponseBody.string()).thenReturn("{\"success\": true}"); Map result = - dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId); + dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId); assertEquals(Map.of("success", true), result); verify(tokenManager, never()) .refreshToken(any(UUID.class), eq(client), any(ObjectMapper.class)); @@ -636,7 +594,7 @@ public void shouldThrowExceptionIfGot413() throws IOException { assertThrows( DestinationMemoryFullException.class, - () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId)); + () -> dtpService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId)); } @Test @@ -660,7 +618,7 @@ public void shouldCallCheckUnprocessableContentIfGot422() assertThrows( NoNasInAccountException.class, - () -> spyService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, requestBody, jobId)); + () -> spyService.sendPostRequest(TestConfigs.TEST_C2_BASE_URL, () -> requestBody, jobId)); } } @@ -715,4 +673,59 @@ public void shouldNotThrowExceptionIfBodyIsInvalidJson() throws CopyExceptionWit dtpService.throwExceptionIfNoQuota(response); } } + + @Nested + public class GetMediaInputStreamWrapperTest { + @Test + public void shouldGetStreamFromJobStore() throws CopyExceptionWithFailureReason, IOException { + // setup + String fetchableUrl = "some_key"; + InputStream inputStream = new ByteArrayInputStream("test data".getBytes()); + long size = 100L; + InputStreamWrapper expected = new InputStreamWrapper(inputStream, size); + when(jobStore.getStream(jobId, fetchableUrl)).thenReturn(expected); + + // act + InputStreamWrapper result = dtpService.getMediaInputStreamWrapper(jobId, fetchableUrl, true); + + // assert + assertEquals(expected, result); + } + + @Test + public void shouldGetStreamFromUrl() throws Exception { + // setup + Path tempFile = Files.createTempFile("test", ".txt"); + tempFile.toFile().deleteOnExit(); + byte[] testData = "test data".getBytes(); + Files.write(tempFile, testData); + String fetchableUrl = tempFile.toUri().toURL().toString(); + + // act + InputStreamWrapper result = dtpService.getMediaInputStreamWrapper(jobId, fetchableUrl, false); + + // assert + assertEquals(testData.length, result.getBytes()); + try (InputStream is = result.getStream()) { + byte[] bytes = new byte[is.available()]; + is.read(bytes); + assertArrayEquals(testData, bytes); + } + } + + @Test + public void shouldThrowExceptionWhenNoSource() { + assertThrows( + IllegalArgumentException.class, + () -> dtpService.getMediaInputStreamWrapper(jobId, null, false)); + } + + @Test + public void shouldThrowExceptionForMalformedUrl() { + String malformedUrl = "this is not a url"; + assertThrows( + IOException.class, + () -> dtpService.getMediaInputStreamWrapper(jobId, malformedUrl, false)); + } + } } diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java index 12f31498c..116d56e36 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/signals/SynologySignalHandlerTest.java @@ -30,13 +30,11 @@ import org.datatransferproject.datatransfer.synology.service.SynologyDTPService; import org.datatransferproject.datatransfer.synology.service.SynologyOAuthTokenManager; import org.datatransferproject.spi.transfer.provider.SignalRequest; -import org.datatransferproject.spi.transfer.types.CopyExceptionWithFailureReason; import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle; import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle.EndReason; import org.datatransferproject.spi.transfer.types.signals.JobLifeCycle.State; import org.datatransferproject.types.common.models.DataVertical; import org.datatransferproject.types.transfer.auth.TokensAndUrlAuthData; -import org.datatransferproject.types.transfer.retry.RetryException; import org.datatransferproject.types.transfer.retry.RetryMapping; import org.datatransferproject.types.transfer.retry.RetryStrategy; import org.datatransferproject.types.transfer.retry.RetryStrategyLibrary; @@ -75,7 +73,7 @@ public void setUp() { } @Test - public void testSendSignal() throws RetryException, CopyExceptionWithFailureReason { + public void testSendSignal() throws Exception { JobLifeCycle jobStatus = JobLifeCycle.builder() .setState(State.ENDED) @@ -100,7 +98,7 @@ public void testSendSignal() throws RetryException, CopyExceptionWithFailureReas } @Test - public void testSendSignalRetry() throws RetryException, CopyExceptionWithFailureReason { + public void testSendSignalRetry() throws Exception { JobLifeCycle jobStatus = JobLifeCycle.builder() .setState(State.ENDED) diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploaderTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploaderTest.java index 42962d3f3..d27def76d 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploaderTest.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/uploader/SynologyUploaderTest.java @@ -115,7 +115,7 @@ public void shouldImportAlbums() throws Exception { } @Test - public void shouldThrowExceptionIfCreateAlbumFails() throws CopyExceptionWithFailureReason { + public void shouldThrowExceptionIfCreateAlbumFails() throws Exception { SynologyUploader spyUploader = Mockito.spy(new SynologyUploader(executor, monitor, synologyDTPService)); List albums = List.of(new MediaAlbum("1", "album1", "desc")); @@ -156,7 +156,7 @@ private Stream provideMediaItems() { } @BeforeEach - public void setUp() throws CopyExceptionWithFailureReason { + public void setUp() throws Exception { lenient() .when(synologyDTPService.addItemToAlbum(any(), any(), any())) .thenReturn(Map.of("success", true)); @@ -190,7 +190,7 @@ public void setUp() throws CopyExceptionWithFailureReason { @ParameterizedTest(name = "shouldImportMediaItemsWhenAlbumBeforeItem [{index}] {0}") @MethodSource("provideMediaItems") public void shouldImportMediaItemsWhenAlbumBeforeItem( - List mediaItems) throws CopyExceptionWithFailureReason { + List mediaItems) throws Exception { SynologyUploader uploader = new SynologyUploader(executor, monitor, synologyDTPService); List albums = List.of(new MediaAlbum(albumId, "album1", "desc")); @@ -268,12 +268,11 @@ public void shouldImportMediaItemsWithCache(List med } } - @ParameterizedTest(name = "shouldThrowExceptionIfCreateMediaItemNotSuccess [{index}] {0}") + @ParameterizedTest(name = "shouldNotThrowExceptionIfCreateMediaItemNotSuccess [{index}] {0}") @MethodSource("provideMediaItems") - public void shouldThrowExceptionIfCreateMediaItemNotSuccess( - List mediaItems) throws CopyExceptionWithFailureReason { - SynologyUploader spyUploader = - Mockito.spy(new SynologyUploader(executor, monitor, synologyDTPService)); + public void shouldNotThrowExceptionIfCreateMediaItemNotSuccess( + List mediaItems) throws Exception { + SynologyUploader uploader = new SynologyUploader(executor, monitor, synologyDTPService); List albums = List.of(new MediaAlbum("1", "album1", "desc")); if (mediaItems.get(0) instanceof PhotoModel) { @@ -282,27 +281,19 @@ public void shouldThrowExceptionIfCreateMediaItemNotSuccess( when(synologyDTPService.createVideo(any(), any())).thenReturn(Map.of("success", false)); } - spyUploader.importAlbums(albums, mockJobId); + uploader.importAlbums(albums, mockJobId); if (mediaItems.get(0) instanceof PhotoModel) { - Exception e = - assertThrows( - CopyExceptionWithFailureReason.class, - () -> spyUploader.importPhotos((List) mediaItems, mockJobId)); - assertTrue(containsMessage(e, "Failed to import item")); + uploader.importPhotos((List) mediaItems, mockJobId); } else if (mediaItems.get(0) instanceof VideoModel) { - Exception e = - assertThrows( - CopyExceptionWithFailureReason.class, - () -> spyUploader.importVideos((List) mediaItems, mockJobId)); - assertTrue(containsMessage(e, "Failed to import item")); + uploader.importVideos((List) mediaItems, mockJobId); } } @ParameterizedTest(name = "shouldThrowExceptionIfAddItemToAlbumFails [{index}] {0}") @MethodSource("provideMediaItems") public void shouldThrowExceptionIfAddItemToAlbumFails( - List mediaItems) throws CopyExceptionWithFailureReason { + List mediaItems) throws Exception { SynologyUploader spyUploader = Mockito.spy(new SynologyUploader(executor, monitor, synologyDTPService)); List albums = List.of(new MediaAlbum("1", "album1", "desc")); @@ -328,31 +319,21 @@ public void shouldThrowExceptionIfAddItemToAlbumFails( } } - @ParameterizedTest(name = "shouldThrowExceptionIfAddItemToAlbumNotSuccess [{index}] {0}") + @ParameterizedTest(name = "shouldNotThrowExceptionIfAddItemToAlbumNotSuccess [{index}] {0}") @MethodSource("provideMediaItems") - public void shouldThrowExceptionIfAddItemToAlbumNotSuccess( - List mediaItems) throws CopyExceptionWithFailureReason { - SynologyUploader spyUploader = - Mockito.spy(new SynologyUploader(executor, monitor, synologyDTPService)); + public void shouldNotThrowExceptionIfAddItemToAlbumNotSuccess( + List mediaItems) throws Exception { + SynologyUploader uploader = new SynologyUploader(executor, monitor, synologyDTPService); List albums = List.of(new MediaAlbum("1", "album1", "desc")); - String expectedMessage = "Unsuccessful"; when(synologyDTPService.addItemToAlbum(any(), any(), any())) .thenReturn(Map.of("success", false)); - spyUploader.importAlbums(albums, mockJobId); + uploader.importAlbums(albums, mockJobId); if (mediaItems.get(0) instanceof PhotoModel) { - Exception e = - assertThrows( - CopyExceptionWithFailureReason.class, - () -> spyUploader.importPhotos((List) mediaItems, mockJobId)); - assertTrue(containsMessage(e, expectedMessage)); + uploader.importPhotos((List) mediaItems, mockJobId); } else if (mediaItems.get(0) instanceof VideoModel) { - Exception e = - assertThrows( - CopyExceptionWithFailureReason.class, - () -> spyUploader.importVideos((List) mediaItems, mockJobId)); - assertTrue(containsMessage(e, expectedMessage)); + uploader.importVideos((List) mediaItems, mockJobId); } } } From d223e20094d23d7e916579ecf404a34ba9f18a26 Mon Sep 17 00:00:00 2001 From: "Chih-Hsien (Simon) Yeh" Date: Thu, 2 Apr 2026 21:40:08 +0800 Subject: [PATCH 04/11] feat(synology): implement chunked upload for large videos (#1488) ## Goal The goal of this change is to provide a more robust upload mechanism for large video files in `SynologyDTPService`. By switching from a single-stream upload to a chunked upload process, we improve reliability for very large transfers and align with the Synology C2 API's preferred method for handling significant media payloads. ## Changes - **Chunked Upload Logic:** Refactored `createVideo` to use a multi-step upload process: - `uploadVideoChunks`: Reads the video stream in 50MB increments and uploads each chunk sequentially to the new `/import/item/chunk` endpoint. - `completeVideoUpload`: Sends a final request to `/import/item/complete` with the total chunk count and metadata (title, description, timestamp) to finalize the file. - **API Configuration:** - Updated `C2Api` and `synology.yaml` to include paths for the new chunk and completion endpoints. - **Client Optimization:** - Updated `configureClient` to force **HTTP/1.1** and increased the default read timeout to 120 seconds to ensure stable long-running connections during chunk transmission. - Simplified `sendPostRequest` by removing the manual timeout override, relying instead on the pre-configured client. - **Memory Efficiency:** Reuses a single byte array buffer for chunking to minimize heap allocations during the transfer of large files. - **Others:** Move file content to the end of multipart ## Testing - **OOM Validation:** Updated `SynologyDTPServiceOOMTest` to verify that a 1GB video is correctly split into multiple chunks and uploaded without exceeding memory limits. - **Functional Tests:** Updated `SynologyDTPServiceTest` to accommodate the new two-step upload flow (Multipart chunks followed by a FormBody completion) and added a specific case `shouldSendMultipleChunksForLargeVideo` to verify correct indexing. --- .../datatransfer/synology/models/C2Api.java | 26 ++ .../synology/service/SynologyDTPService.java | 189 +++++++----- .../src/main/resources/config/synology.yaml | 2 + .../service/SynologyDTPServiceOOMTest.java | 40 ++- .../service/SynologyDTPServiceTest.java | 288 +++++++++++++----- .../synology/utils/TestConfigs.java | 4 + 6 files changed, 385 insertions(+), 164 deletions(-) diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java index 1cc3ef799..bdc5dd397 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/models/C2Api.java @@ -27,6 +27,8 @@ public class C2Api extends ServiceConfig.Service { private final String createAlbum; private final String uploadItem; + private final String chunkUploadItem; + private final String completeUploadItem; private final String addItemToAlbum; private final String signalJob; @@ -38,6 +40,8 @@ public C2Api(@JsonProperty("baseUrl") String baseUrl, @JsonProperty("apiPath") A this.createAlbum = UrlUtils.join(baseUrl, apiPath.getCreateAlbumPath()); this.uploadItem = UrlUtils.join(baseUrl, apiPath.getUploadItemPath()); + this.chunkUploadItem = UrlUtils.join(baseUrl, apiPath.getChunkUploadItemPath()); + this.completeUploadItem = UrlUtils.join(baseUrl, apiPath.getCompleteUploadItemPath()); this.addItemToAlbum = UrlUtils.join(baseUrl, apiPath.getAddItemToAlbumPath()); this.signalJob = UrlUtils.join(baseUrl, apiPath.getSignalJobPath()); } @@ -54,6 +58,14 @@ public String getUploadItem() { return uploadItem; } + public String getChunkUploadItem() { + return chunkUploadItem; + } + + public String getCompleteUploadItem() { + return completeUploadItem; + } + public String getAddItemToAlbum() { return addItemToAlbum; } @@ -65,6 +77,8 @@ public String getSignalJob() { public static class ApiPath { private final String createAlbumPath; private final String uploadItemPath; + private final String chunkUploadItemPath; + private final String completeUploadItemPath; private final String addItemToAlbumPath; private final String signalJobPath; @@ -72,10 +86,14 @@ public static class ApiPath { public ApiPath( @JsonProperty("createAlbum") String createAlbumPath, @JsonProperty("uploadItem") String uploadItemPath, + @JsonProperty("chunkUploadItem") String chunkUploadItemPath, + @JsonProperty("completeUploadItem") String completeUploadItemPath, @JsonProperty("addItemToAlbum") String addItemToAlbumPath, @JsonProperty("signalJob") String signalJobPath) { this.createAlbumPath = createAlbumPath; this.uploadItemPath = uploadItemPath; + this.chunkUploadItemPath = chunkUploadItemPath; + this.completeUploadItemPath = completeUploadItemPath; this.addItemToAlbumPath = addItemToAlbumPath; this.signalJobPath = signalJobPath; } @@ -88,6 +106,14 @@ public String getUploadItemPath() { return uploadItemPath; } + public String getChunkUploadItemPath() { + return chunkUploadItemPath; + } + + public String getCompleteUploadItemPath() { + return completeUploadItemPath; + } + public String getAddItemToAlbumPath() { return addItemToAlbumPath; } diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java index d69202585..78b7ac460 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java @@ -21,8 +21,10 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Strings; import java.io.IOException; +import java.io.InputStream; import java.net.URL; import java.net.URLConnection; +import java.util.Arrays; import java.util.Date; import java.util.Map; import java.util.UUID; @@ -93,7 +95,16 @@ public SynologyDTPService( this.objectMapper = new ObjectMapper(); this.jobStore = jobStore; this.tokenManager = tokenManager; - this.client = client; + this.client = configureClient(client); + } + + @VisibleForTesting + protected OkHttpClient configureClient(OkHttpClient client) { + return client + .newBuilder() + .protocols(Arrays.asList(okhttp3.Protocol.HTTP_1_1)) + .readTimeout(120, java.util.concurrent.TimeUnit.SECONDS) + .build(); } /** @@ -190,7 +201,6 @@ public void writeTo(BufferedSink sink) throws IOException { MultipartBody.Builder builder = new MultipartBody.Builder() .setType(MultipartBody.FORM) - .addFormDataPart("file", photo.getTitle(), fileBody) .addFormDataPart("item_id", photo.getDataId()) .addFormDataPart("title", photo.getTitle()) .addFormDataPart("job_id", jobId.toString()) @@ -206,6 +216,8 @@ public void writeTo(BufferedSink sink) throws IOException { builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); } + builder.addFormDataPart("file_size", String.valueOf(inputStreamWrapper.getBytes())); + builder.addFormDataPart("file", photo.getTitle(), fileBody); return builder.build(); }; @@ -237,70 +249,110 @@ public Map createVideo(VideoModel video, UUID jobId) video.getDataId(), video.getName()), jobId); - RequestBodyGenerator bodyGenerator = - () -> { - // Due to InputStream may not repeatable, we need to open it inside the generator function - // to make sure it can be read when retrying. - InputStreamWrapper inputStreamWrapper = - getMediaInputStreamWrapper(jobId, video.getFetchableUrl(), video.isInTempStore()); + InputStreamWrapper inputStreamWrapper = + getMediaInputStreamWrapper(jobId, video.getFetchableUrl(), video.isInTempStore()); - RequestBody fileBody = - new RequestBody() { - private boolean isConsumed = false; + int actualChunkCount; + try (InputStream inputStream = inputStreamWrapper.getStream()) { + actualChunkCount = uploadVideoChunks(video, jobId, inputStream); + } - @Override - public MediaType contentType() { - return MediaType.parse(video.getMimeType()); - } + Map responseData = completeVideoUpload(video, jobId, actualChunkCount); - @Override - public long contentLength() { - return inputStreamWrapper.getBytes(); - } + monitor.info( + () -> + String.format( + "[SynologyImporter] video created successfully, name: [%s].", video.getName()), + jobId); + return responseData; + } - @Override - public void writeTo(BufferedSink sink) throws IOException { - if (isConsumed) { - throw new IOException("InputStream has already been consumed"); - } - isConsumed = true; - try (Source source = Okio.source(inputStreamWrapper.getStream())) { - sink.writeAll(source); - } - } - }; + private int uploadVideoChunks(VideoModel video, UUID jobId, InputStream inputStream) + throws CopyExceptionWithFailureReason, IOException { + final int CHUNK_SIZE = 50 * 1024 * 1024; // 50MB + int actualChunkCount = 0; + + // reuse the same byte array for each chunk to reduce memory usage, + // since we are uploading in sequential, there is no concurrency issue + byte[] chunkData = new byte[CHUNK_SIZE]; + while (true) { + int bytesRead = inputStream.readNBytes(chunkData, 0, CHUNK_SIZE); + if (bytesRead == 0) { + break; + } - MultipartBody.Builder builder = - new MultipartBody.Builder() - .setType(MultipartBody.FORM) - .addFormDataPart("file", video.getName(), fileBody) - .addFormDataPart("item_id", video.getDataId()) - .addFormDataPart("title", video.getName()) - .addFormDataPart("job_id", jobId.toString()) - .addFormDataPart("service", exportingService); + RequestBody fileBody = + new RequestBody() { + @Override + public MediaType contentType() { + return MediaType.parse(video.getMimeType()); + } - String imageDescription = video.getDescription(); - if (!Strings.isNullOrEmpty(imageDescription)) { - builder.addFormDataPart("description", imageDescription); - } - Date videoUploadedTime = video.getUploadedTime(); - if (videoUploadedTime != null) { - long timestampInSeconds = videoUploadedTime.getTime() / 1000; - builder.addFormDataPart("uploaded_time", String.valueOf(timestampInSeconds)); - } + @Override + public long contentLength() { + return bytesRead; + } - return builder.build(); - }; + @Override + public void writeTo(BufferedSink sink) throws IOException { + // write the actual bytes read for the last chunk, which may be smaller than + // CHUNK_SIZE + sink.write(chunkData, 0, bytesRead); + } + }; + + MultipartBody.Builder builder = + new MultipartBody.Builder() + .setType(MultipartBody.FORM) + .addFormDataPart("item_id", video.getDataId()) + .addFormDataPart("index", String.valueOf(actualChunkCount)) + .addFormDataPart("job_id", jobId.toString()) + .addFormDataPart("service", exportingService) + .addFormDataPart("file_size", String.valueOf(bytesRead)) + .addFormDataPart("file", video.getName(), fileBody); + + RequestBody requestBody = builder.build(); + + final String chunkInfo = + String.format( + "[SynologyImporter] uploading video chunk, video name: [%s], chunk index: [%d], chunk" + + " size: [%d].", + video.getName(), actualChunkCount, bytesRead); + monitor.info(() -> chunkInfo, jobId); + sendPostRequest(c2Api.getChunkUploadItem(), () -> requestBody, jobId); + + actualChunkCount++; + } + return actualChunkCount; + } + private Map completeVideoUpload(VideoModel video, UUID jobId, int totalChunks) + throws CopyExceptionWithFailureReason, IOException { + FormBody.Builder builder = + new FormBody.Builder() + .add("item_id", video.getDataId()) + .add("title", video.getName()) + .add("total_chunks", String.valueOf(totalChunks)) + .add("job_id", jobId.toString()) + .add("service", exportingService); + + if (!Strings.isNullOrEmpty(video.getDescription())) { + builder.add("description", video.getDescription()); + } + if (video.getUploadedTime() != null) { + builder.add("uploaded_time", String.valueOf(video.getUploadedTime().getTime() / 1000)); + } + RequestBody requestBody = builder.build(); + + final String completeInfo = + String.format( + "[SynologyImporter] completing video upload, video name: [%s], total chunks: [%d].", + video.getName(), totalChunks); + monitor.info(() -> completeInfo, jobId); @SuppressWarnings("unchecked") Map responseData = (Map) - sendPostRequest(c2Api.getUploadItem(), bodyGenerator, jobId, 300).get("data"); - monitor.info( - () -> - String.format( - "[SynologyImporter] video created successfully, name: [%s].", video.getName()), - jobId); + sendPostRequest(c2Api.getCompleteUploadItem(), () -> requestBody, jobId).get("data"); return responseData; } @@ -380,13 +432,6 @@ protected void throwExceptionIfNoQuota(Response response) throws CopyExceptionWi } } - @VisibleForTesting - protected Map sendPostRequest( - String url, RequestBodyGenerator bodyGenerator, UUID jobId) - throws CopyExceptionWithFailureReason, IOException { - return sendPostRequest(url, bodyGenerator, jobId, -1); - } - /* * @param url the URL to send the POST request to * @param bodyGenerator a generator function that produces the request body, it will be called for each retry attempt @@ -395,29 +440,23 @@ protected Map sendPostRequest( */ @VisibleForTesting protected Map sendPostRequest( - String url, RequestBodyGenerator bodyGenerator, UUID jobId, int timeoutInSeconds) + String url, RequestBodyGenerator bodyGenerator, UUID jobId) throws CopyExceptionWithFailureReason, IOException { boolean triedRefreshToken = false; Exception lastException = null; for (int retry = retryConfig.getMaxAttempts(); retry > 0; retry--) { + final String methodInfo = + String.format( + "[SynologyImporter] Sending POST request to url: [%s], attempts left: [%d]", + url, retry); + monitor.info(() -> methodInfo, jobId); Response response = null; try { RequestBody body = bodyGenerator.get(); Request.Builder requestBuilder = new Request.Builder().url(url).post(body); requestBuilder.header("Authorization", "Bearer " + tokenManager.getAccessToken(jobId)); - if (timeoutInSeconds < 0) { - // Use default timeout of OkHttpClient - response = client.newCall(requestBuilder.build()).execute(); - } else { - response = - client - .newBuilder() - .readTimeout(timeoutInSeconds, java.util.concurrent.TimeUnit.SECONDS) - .build() - .newCall(requestBuilder.build()) - .execute(); - } + response = client.newCall(requestBuilder.build()).execute(); if (!response.isSuccessful()) { int code = response.code(); if (code == 401) { diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml b/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml index aefd9d4f1..90a6818c7 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/resources/config/synology.yaml @@ -10,5 +10,7 @@ serviceConfig: apiPath: createAlbum: "/import/album" uploadItem: "/import/item" + chunkUploadItem: "/import/item/chunk" + completeUploadItem: "/import/item/complete" addItemToAlbum: "/import/album/item" signalJob: "/import/job/signal" diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java index 48b5510cd..90db54c11 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceOOMTest.java @@ -19,13 +19,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.when; import java.io.IOException; import java.io.InputStream; +import java.util.List; import java.util.Map; import java.util.UUID; import okhttp3.OkHttpClient; @@ -136,7 +136,12 @@ public void setUp() throws InvalidTokenException { .thenReturn(TestConfigs.createServiceConfigJson()); dtpService = new SynologyDTPService( - monitor, transferServiceConfig, exportingService, jobStore, tokenManager, client); + monitor, transferServiceConfig, exportingService, jobStore, tokenManager, client) { + @Override + protected OkHttpClient configureClient(OkHttpClient client) { + return client; + } + }; } @Test @@ -198,30 +203,37 @@ public void createVideo_withLargeFile_shouldStreamData() throws Exception { SynologyDTPService spyService = Mockito.spy(dtpService); doReturn(Map.of("success", true, "data", Map.of("item_id", "video-id"))) .when(spyService) - .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); // act spyService.createVideo(video, jobId); - // Get the generated request body - RequestBody requestBody = requestBodyCaptor.getValue().get(); + // Get the generated request bodies + List generators = requestBodyCaptor.getAllValues(); // assert - // The content length is not asserted to be equal to 1GB because it includes multipart - // boundaries and headers, - // making it larger than the raw file size. The main purpose of this test is to ensure that - // the file is streamed without causing an OutOfMemoryError. - assertTrue(requestBody.contentLength() > oneGB); + // The total content length across all chunks should be at least 1GB + long totalContentLength = 0; - // Simulate writing the body to a sink that discards the data. + // Simulate writing the bodies to a sink that discards the data. // This will throw OutOfMemoryError if the whole stream is loaded into memory. okio.BufferedSink discardingSink = Okio.buffer(new BlackholeSink()); - // The MultipartBody will try to read the stream to write it. - // If it buffers the whole 1GB in memory, this will OOM. - requestBody.writeTo(discardingSink); + for (SynologyDTPService.RequestBodyGenerator generator : generators) { + RequestBody requestBody = generator.get(); + long length = requestBody.contentLength(); + // Only count chunks which are significantly large to avoid counting the "complete" request + if (length > 1024 * 1024) { + totalContentLength += length; + } + requestBody.writeTo(discardingSink); + } discardingSink.flush(); + // The total content length should be at least 1GB (actually more due to multipart boundaries + // and headers) + assertTrue(totalContentLength > oneGB); + // If we reach here, it means we streamed the data without OOM. } } diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java index cb6757905..851116589 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPServiceTest.java @@ -25,7 +25,6 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.any; -import static org.mockito.Mockito.anyInt; import static org.mockito.Mockito.anyString; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; @@ -109,7 +108,12 @@ public void setUp() throws InvalidTokenException { dtpService = new SynologyDTPService( - monitor, transferServiceConfig, exportingService, jobStore, tokenManager, client); + monitor, transferServiceConfig, exportingService, jobStore, tokenManager, client) { + @Override + protected OkHttpClient configureClient(OkHttpClient client) { + return client; + } + }; } @Nested @@ -262,45 +266,103 @@ public void shouldSendPostRequestWithCorrectFormBodyWithDescriptionAndUploadedTi .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createPhoto((PhotoModel) item, jobId); verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + + assertEquals(result.get("item_id"), itemId); + assertEquals(result.get("success"), null); + + RequestBody capturedBody = requestBodyCaptor.getValue().get(); + MultipartBody multipartBody = (MultipartBody) capturedBody; + + Map multipartFormAnswer = + Map.of( + "job_id", jobId.toString(), + "service", exportingService, + "item_id", itemId, + "title", itemName, + "description", description, + "uploaded_time", String.valueOf(uploadedTimestampInSeconds)); + + for (MultipartBody.Part part : multipartBody.parts()) { + String partName = + part.headers() + .get("Content-Disposition") + .split(";")[1] + .split("=")[1] + .replace("\"", ""); + Buffer buffer = new Buffer(); + part.body().writeTo(buffer); + if (partName.equals("file")) { + byte[] partBytes = buffer.readByteArray(); + + assertArrayEquals(mockImage, partBytes); + String fileName = item.getName(); + String contentDisposition = part.headers().get("Content-Disposition"); + assertTrue(contentDisposition.contains("filename=\"" + fileName + "\"")); + } else if (multipartFormAnswer.containsKey(partName)) { + String partValue = buffer.readUtf8(); + assertEquals(multipartFormAnswer.get(partName), partValue); + } + } } else if (item instanceof VideoModel) { doReturn(Map.of("success", true, "data", dataMap)) .when(spyService) - .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createVideo((VideoModel) item, jobId); - verify(spyService) - .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); - } - - assertEquals(result.get("item_id"), itemId); - assertEquals(result.get("success"), null); - - RequestBody capturedBody = requestBodyCaptor.getValue().get(); - MultipartBody multipartBody = (MultipartBody) capturedBody; - - Map multipartFormAnswer = - Map.of( - "job_id", jobId.toString(), - "service", exportingService, - "item_id", itemId, - "title", itemName, - "description", description, - "uploaded_time", String.valueOf(uploadedTimestampInSeconds)); + verify(spyService, times(2)) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + + assertEquals(result.get("item_id"), itemId); + assertEquals(result.get("success"), null); + + java.util.List capturedGenerators = requestBodyCaptor.getAllValues(); + assertEquals(2, capturedGenerators.size()); + + // Check first call (chunk) + RequestBody chunkBody = capturedGenerators.get(0).get(); + assertTrue(chunkBody instanceof MultipartBody); + MultipartBody chunkMultipart = (MultipartBody) chunkBody; + Map chunkAnswer = + Map.of( + "job_id", + jobId.toString(), + "service", + exportingService, + "item_id", + itemId, + "index", + "0"); + for (MultipartBody.Part part : chunkMultipart.parts()) { + String partName = + part.headers() + .get("Content-Disposition") + .split(";")[1] + .split("=")[1] + .replace("\"", ""); + Buffer buffer = new Buffer(); + part.body().writeTo(buffer); + if (partName.equals("file")) { + byte[] partBytes = buffer.readByteArray(); + assertArrayEquals(mockImage, partBytes); + } else if (chunkAnswer.containsKey(partName)) { + assertEquals(chunkAnswer.get(partName), buffer.readUtf8()); + } + } - for (MultipartBody.Part part : multipartBody.parts()) { - String partName = - part.headers().get("Content-Disposition").split(";")[1].split("=")[1].replace("\"", ""); - Buffer buffer = new Buffer(); - part.body().writeTo(buffer); - if (partName.equals("file")) { - byte[] partBytes = buffer.readByteArray(); - - assertArrayEquals(mockImage, partBytes); - String fileName = item.getName(); - String contentDisposition = part.headers().get("Content-Disposition"); - assertTrue(contentDisposition.contains("filename=\"" + fileName + "\"")); - } else if (multipartFormAnswer.containsKey(partName)) { - String partValue = buffer.readUtf8(); - assertEquals(multipartFormAnswer.get(partName), partValue); + // Check second call (complete) + RequestBody completeBody = capturedGenerators.get(1).get(); + assertTrue(completeBody instanceof FormBody); + FormBody completeForm = (FormBody) completeBody; + Map completeAnswer = + Map.of( + "job_id", jobId.toString(), + "service", exportingService, + "item_id", itemId, + "title", itemName, + "total_chunks", "1", + "description", description, + "uploaded_time", String.valueOf(uploadedTimestampInSeconds)); + for (int i = 0; i < completeForm.size(); i++) { + assertEquals(completeAnswer.get(completeForm.name(i)), completeForm.value(i)); } } } @@ -327,45 +389,65 @@ public void shouldSendPostRequestWithCorrectFormBodyWithoutDescription(Downloada .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createPhoto((PhotoModel) item, jobId); verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + + assertEquals(result.get("item_id"), itemId); + assertEquals(result.get("success"), null); + + RequestBody capturedBody = requestBodyCaptor.getValue().get(); + MultipartBody multipartBody = (MultipartBody) capturedBody; + + Map multipartFormAnswer = + Map.of( + "job_id", jobId.toString(), + "service", exportingService, + "item_id", itemId, + "title", itemName); + + for (MultipartBody.Part part : multipartBody.parts()) { + String partName = + part.headers() + .get("Content-Disposition") + .split(";")[1] + .split("=")[1] + .replace("\"", ""); + Buffer buffer = new Buffer(); + part.body().writeTo(buffer); + assertNotEquals("description", partName); + if (partName.equals("file")) { + byte[] partBytes = buffer.readByteArray(); + + assertArrayEquals(mockImage, partBytes); + String fileName = item.getName(); + String contentDisposition = part.headers().get("Content-Disposition"); + assertTrue(contentDisposition.contains("filename=\"" + fileName + "\"")); + } else if (multipartFormAnswer.containsKey(partName)) { + String partValue = buffer.readUtf8(); + assertEquals(multipartFormAnswer.get(partName), partValue); + } + } } else if (item instanceof VideoModel) { doReturn(Map.of("success", true, "data", dataMap)) .when(spyService) - .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); result = spyService.createVideo((VideoModel) item, jobId); - verify(spyService) - .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); - } + verify(spyService, times(2)) + .sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); - assertEquals(result.get("item_id"), itemId); - assertEquals(result.get("success"), null); + assertEquals(result.get("item_id"), itemId); + assertEquals(result.get("success"), null); - RequestBody capturedBody = requestBodyCaptor.getValue().get(); - MultipartBody multipartBody = (MultipartBody) capturedBody; + java.util.List capturedGenerators = requestBodyCaptor.getAllValues(); + assertEquals(2, capturedGenerators.size()); - Map multipartFormAnswer = - Map.of( - "job_id", jobId.toString(), - "service", exportingService, - "item_id", itemId, - "title", itemName); + // Check first call (chunk) + RequestBody chunkBody = capturedGenerators.get(0).get(); + assertTrue(chunkBody instanceof MultipartBody); - for (MultipartBody.Part part : multipartBody.parts()) { - String partName = - part.headers().get("Content-Disposition").split(";")[1].split("=")[1].replace("\"", ""); - Buffer buffer = new Buffer(); - part.body().writeTo(buffer); - assertNotEquals("description", partName); - if (partName.equals("file")) { - byte[] partBytes = buffer.readByteArray(); - - assertArrayEquals(mockImage, partBytes); - String fileName = item.getName(); - String contentDisposition = part.headers().get("Content-Disposition"); - assertTrue(contentDisposition.contains("filename=\"" + fileName + "\"")); - } else if (multipartFormAnswer.containsKey(partName)) { - String partValue = buffer.readUtf8(); - assertEquals(multipartFormAnswer.get(partName), partValue); - } + // Check second call (complete) + RequestBody completeBody = capturedGenerators.get(1).get(); + assertTrue(completeBody instanceof FormBody); + FormBody completeForm = (FormBody) completeBody; + assertEquals("1", completeForm.value(2)); // total_chunks } } @@ -373,6 +455,11 @@ public void shouldSendPostRequestWithCorrectFormBodyWithoutDescription(Downloada @MethodSource("provideMediaObjectsInTempStore") public void shouldThrowExceptionIfInputStreamIsConsumed(DownloadableFile item) throws IOException, CopyExceptionWithFailureReason { + if (item instanceof VideoModel) { + // VideoModel uses chunked upload with buffered byte arrays for retries, + // so RequestBody can be written multiple times. + return; + } byte[] mockImage = new byte[] {1, 2, 3}; InputStream mockInputStream = new ByteArrayInputStream(mockImage); InputStreamWrapper streamWrapper = mock(InputStreamWrapper.class); @@ -387,13 +474,6 @@ public void shouldThrowExceptionIfInputStreamIsConsumed(DownloadableFile item) .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); spyService.createPhoto((PhotoModel) item, jobId); verify(spyService).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); - } else if (item instanceof VideoModel) { - doReturn(Map.of("success", true, "data", dataMap)) - .when(spyService) - .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); - spyService.createVideo((VideoModel) item, jobId); - verify(spyService) - .sendPostRequest(anyString(), requestBodyCaptor.capture(), any(), anyInt()); } RequestBody capturedBody = requestBodyCaptor.getValue().get(); @@ -403,6 +483,62 @@ public void shouldThrowExceptionIfInputStreamIsConsumed(DownloadableFile item) assertTrue(exception.getMessage().contains("InputStream has already been consumed")); } + @Test + public void shouldSendMultipleChunksForLargeVideo() + throws IOException, CopyExceptionWithFailureReason { + // 50MB + 1 byte to trigger 2 chunks + int firstChunkSize = 50 * 1024 * 1024; + byte[] mockImage = new byte[firstChunkSize + 1]; + + InputStream mockInputStream = new ByteArrayInputStream(mockImage); + InputStreamWrapper streamWrapper = mock(InputStreamWrapper.class); + SynologyDTPService spyService = Mockito.spy(dtpService); + Map dataMap = Map.of("item_id", itemId); + VideoModel video = + new VideoModel(itemName, fetchUrl, description, "format", itemId, null, true, null); + + when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); + when(streamWrapper.getStream()).thenReturn(mockInputStream); + doReturn(Map.of("success", true, "data", dataMap)) + .when(spyService) + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); + + spyService.createVideo(video, jobId); + + // Should call sendPostRequest 3 times: chunk 0, chunk 1, complete + verify(spyService, times(3)).sendPostRequest(anyString(), requestBodyCaptor.capture(), any()); + + java.util.List capturedGenerators = requestBodyCaptor.getAllValues(); + assertEquals(3, capturedGenerators.size()); + + // Check chunk 0 + RequestBody chunk0Body = capturedGenerators.get(0).get(); + assertEquals("0", getMultipartPartValue((MultipartBody) chunk0Body, "index")); + + // Check chunk 1 + RequestBody chunk1Body = capturedGenerators.get(1).get(); + assertEquals("1", getMultipartPartValue((MultipartBody) chunk1Body, "index")); + + // Check complete + RequestBody completeBody = capturedGenerators.get(2).get(); + FormBody completeForm = (FormBody) completeBody; + assertEquals("2", completeForm.value(2)); // total_chunks + } + + private String getMultipartPartValue(MultipartBody multipartBody, String partName) + throws IOException { + for (MultipartBody.Part part : multipartBody.parts()) { + String name = + part.headers().get("Content-Disposition").split(";")[1].split("=")[1].replace("\"", ""); + if (name.equals(partName)) { + Buffer buffer = new Buffer(); + part.body().writeTo(buffer); + return buffer.readUtf8(); + } + } + return null; + } + @ParameterizedTest(name = "shouldThrowExceptionIfSendPostRequestFailed [{index}] {0}") @MethodSource("provideMediaObjectsInTempStore") public void shouldThrowExceptionIfSendPostRequestFailed(DownloadableFile item) @@ -421,9 +557,11 @@ public void shouldThrowExceptionIfSendPostRequestFailed(DownloadableFile item) () -> spyService.createPhoto((PhotoModel) item, jobId), "MockException"); } else if (item instanceof VideoModel) { + when(jobStore.getStream(jobId, fetchUrl)).thenReturn(streamWrapper); + when(streamWrapper.getStream()).thenReturn(mockInputStream); doThrow(new UploadErrorException("MockException", null)) .when(spyService) - .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any(), anyInt()); + .sendPostRequest(anyString(), any(RequestBodyGenerator.class), any()); assertThrows( UploadErrorException.class, () -> spyService.createVideo((VideoModel) item, jobId), diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java index fa2280a49..1eebaa74e 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/test/java/org/datatransferproject/datatransfer/synology/utils/TestConfigs.java @@ -29,6 +29,8 @@ public class TestConfigs { public static final String TEST_C2_BASE_URL = "https://fake.url"; public static final String TEST_CREATE_ALBUM_PATH = "/create"; public static final String TEST_UPLOAD_ITEM_PATH = "/upload"; + public static final String TEST_CHUNK_UPLOAD_ITEM_PATH = "/chunk"; + public static final String TEST_COMPLETE_UPLOAD_ITEM_PATH = "/complete"; public static final String TEST_ADD_ITEM_TO_ALBUM_PATH = "/add"; public static final String TEST_SIGNAL_JOB_PATH = "/signal"; public static final int TEST_MAX_ATTEMPTS = 5; @@ -38,6 +40,8 @@ public static ServiceConfig createServiceConfig() { new C2Api.ApiPath( TEST_CREATE_ALBUM_PATH, TEST_UPLOAD_ITEM_PATH, + TEST_CHUNK_UPLOAD_ITEM_PATH, + TEST_COMPLETE_UPLOAD_ITEM_PATH, TEST_ADD_ITEM_TO_ALBUM_PATH, TEST_SIGNAL_JOB_PATH); From 34f48944b75c4cbaf66b8139b3530b917774a5ca Mon Sep 17 00:00:00 2001 From: Aman Pratik Date: Tue, 7 Apr 2026 04:36:55 +0530 Subject: [PATCH 05/11] fix: AppleSignalInterface token refresh (#1487) Fixing `AppleSignalInterface:: sendPostRequest` to accomodate Access-Token refresh POST calls. Co-authored-by: aman-pratik Co-authored-by: Sundeep Paruvu --- .../apple/signals/AppleSignalInterface.java | 27 ++++++++++--------- 1 file changed, 14 insertions(+), 13 deletions(-) diff --git a/extensions/data-transfer/portability-data-transfer-apple/src/main/java/org/datatransferproject/datatransfer/apple/signals/AppleSignalInterface.java b/extensions/data-transfer/portability-data-transfer-apple/src/main/java/org/datatransferproject/datatransfer/apple/signals/AppleSignalInterface.java index b02d2bd40..2b945f661 100644 --- a/extensions/data-transfer/portability-data-transfer-apple/src/main/java/org/datatransferproject/datatransfer/apple/signals/AppleSignalInterface.java +++ b/extensions/data-transfer/portability-data-transfer-apple/src/main/java/org/datatransferproject/datatransfer/apple/signals/AppleSignalInterface.java @@ -17,30 +17,23 @@ package org.datatransferproject.datatransfer.apple.signals; import com.fasterxml.jackson.databind.ObjectMapper; -import com.google.common.annotations.VisibleForTesting; import java.io.IOException; -import java.io.Serializable; import java.net.HttpURLConnection; import java.net.URL; import java.nio.charset.StandardCharsets; -import java.util.Map; import java.util.UUID; import org.apache.commons.io.IOUtils; -import org.apache.commons.lang3.SerializationUtils; import org.datatransferproject.api.launcher.Monitor; import org.datatransferproject.datatransfer.apple.AppleBaseInterface; -import org.datatransferproject.datatransfer.apple.AppleInterfaceFactory; import org.datatransferproject.datatransfer.apple.constants.AuditKeys; import org.datatransferproject.datatransfer.apple.constants.Headers; import org.datatransferproject.spi.transfer.provider.SignalRequest; -import org.datatransferproject.spi.transfer.types.CopyException; import org.datatransferproject.spi.transfer.types.CopyExceptionWithFailureReason; import org.datatransferproject.spi.transfer.types.PermissionDeniedException; import org.datatransferproject.spi.transfer.types.UnconfirmedUserException; import org.datatransferproject.transfer.JobMetadata; import org.datatransferproject.types.transfer.auth.AppCredentials; import org.datatransferproject.types.transfer.auth.TokensAndUrlAuthData; -import org.datatransferproject.types.transfer.retry.RetryStrategyLibrary; import org.jetbrains.annotations.NotNull; /** @@ -48,6 +41,8 @@ */ public class AppleSignalInterface implements AppleBaseInterface { private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + private static final String MEDIA_TYPE_JSON = "application/json"; + private static final String MEDIA_TYPE_FORM = "application/x-www-form-urlencoded"; protected String baseUrl; protected AppCredentials appCredentials; @@ -89,8 +84,12 @@ public String sendPostRequest(@NotNull String url, @NotNull final byte[] request con.setRequestMethod("POST"); con.setRequestProperty(Headers.AUTHORIZATION.getValue(), authData.getAccessToken()); con.setRequestProperty(Headers.CORRELATION_ID.getValue(), correlationId); - con.setRequestProperty(Headers.CONTENT_TYPE.getValue(), "application/json"); - con.setRequestProperty(Headers.ACCEPT.getValue(), "application/json"); + con.setRequestProperty(Headers.ACCEPT.getValue(), MEDIA_TYPE_JSON); + if (url.equals(authData.getTokenServerEncodedUrl())) { + con.setRequestProperty(Headers.CONTENT_TYPE.getValue(), MEDIA_TYPE_FORM); + } else { + con.setRequestProperty(Headers.CONTENT_TYPE.getValue(), MEDIA_TYPE_JSON); + } IOUtils.write(requestData, con.getOutputStream()); responseString = IOUtils.toString(con.getInputStream(), StandardCharsets.ISO_8859_1); @@ -106,15 +105,17 @@ public String sendPostRequest(@NotNull String url, @NotNull final byte[] request AuditKeys.error, e.getMessage(), AuditKeys.errorCode, - con.getResponseCode(), + (con != null ? con.getResponseCode() : -1), e); convertAndThrowException(e, con); } finally { - con.disconnect(); + if (con != null) { + con.disconnect(); + } } - return responseString; } - return null; + + return responseString; } public byte[] sendSignal(@NotNull final SignalRequest signalRequest) From f20db878c8ffc90a19c9d4d2e0bb8009e1703711 Mon Sep 17 00:00:00 2001 From: "Chih-Hsien (Simon) Yeh" Date: Tue, 7 Apr 2026 14:41:39 +0800 Subject: [PATCH 06/11] fix(synology): replace readNBytes with Guava ByteStreams.read (#1489) ## Goal The goal of this change is to improve compatibility of chunked reading in by replacing the Java 9+ method with Guava's , ensuring the code works correctly in older Java environments. ## Changes - **Dependency Update:** Added `com.google.guava:guava` to the Synology extension's `build.gradle`. - **Code Refactor:** Updated `SynologyDTPService` to use `ByteStreams.read` when reading video chunks. Co-authored-by: emma --- .../portability-data-transfer-synology/build.gradle | 1 + .../datatransfer/synology/service/SynologyDTPService.java | 3 ++- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/extensions/data-transfer/portability-data-transfer-synology/build.gradle b/extensions/data-transfer/portability-data-transfer-synology/build.gradle index bc1808c6f..f47862dea 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/build.gradle +++ b/extensions/data-transfer/portability-data-transfer-synology/build.gradle @@ -25,6 +25,7 @@ dependencies { compile project(':portability-spi-cloud') compile project(':portability-transfer') compile "com.squareup.okhttp3:okhttp:${okHttpVersion}" + compile "com.google.guava:guava:${guavaVersion}" } configurePublication(project) diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java index 78b7ac460..bceb1fb28 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java @@ -20,6 +20,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Strings; +import com.google.common.io.ByteStreams; import java.io.IOException; import java.io.InputStream; import java.net.URL; @@ -276,7 +277,7 @@ private int uploadVideoChunks(VideoModel video, UUID jobId, InputStream inputStr // since we are uploading in sequential, there is no concurrency issue byte[] chunkData = new byte[CHUNK_SIZE]; while (true) { - int bytesRead = inputStream.readNBytes(chunkData, 0, CHUNK_SIZE); + int bytesRead = ByteStreams.read(inputStream, chunkData, 0, CHUNK_SIZE); if (bytesRead == 0) { break; } From 5689feb497d21497723a42d98d7629bdc9da51a2 Mon Sep 17 00:00:00 2001 From: Lisa Dusseault Date: Tue, 7 Apr 2026 11:38:08 -0700 Subject: [PATCH 07/11] chore: Cleanup: Removes old log4j values and file left over (#1486) fixes #1031, which has some of the investigation confirming this was really left over and now unused --- .../src/main/resources/log4j.properties | 21 ------------------- gradle.properties | 1 - 2 files changed, 22 deletions(-) delete mode 100644 distributions/demo-server/src/main/resources/log4j.properties diff --git a/distributions/demo-server/src/main/resources/log4j.properties b/distributions/demo-server/src/main/resources/log4j.properties deleted file mode 100644 index 4c8dce47f..000000000 --- a/distributions/demo-server/src/main/resources/log4j.properties +++ /dev/null @@ -1,21 +0,0 @@ -# -# Copyright 2018 The Data Transfer Project Authors. -# -# Licensed under the Apache License, Version 2.0 (the "License"); -# you may not use this file except in compliance with the License. -# You may obtain a copy of the License at -# -# https://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# - -# SLF4J logging configuration - -log4j.rootLogger=DEBUG, consoleAppender -log4j.appender.consoleAppender=org.apache.log4j.ConsoleAppender -log4j.appender.consoleAppender.layout=org.datatransferproject.logging.EncryptingLayout diff --git a/gradle.properties b/gradle.properties index d7ed9cdd0..95882afc9 100644 --- a/gradle.properties +++ b/gradle.properties @@ -38,7 +38,6 @@ okHttpVersion=3.9.1 azureKeyVaultVersion=1.0.0 javaDockerContainer=openjdk:11 jerseyVersion=2.26 -log4jVersion=2.17.0 apacheHttpVersion=4.5.6 # ossrhUsername= From 57744c76aecf307321ee902ba613404a88f8ec60 Mon Sep 17 00:00:00 2001 From: "Chih-Hsien (Simon) Yeh" Date: Thu, 30 Apr 2026 14:56:11 +0800 Subject: [PATCH 08/11] feat(synology): optimize video upload with chunk preloading (#1490) Improve video upload efficiency by implementing an asynchronous producer-consumer model for chunked uploads. This minimizes idle time between network requests and increases overall throughput for large video files. ## Goal The goal of this change is to optimize video uploads to Synology by pre-fetching the next data chunk while the current one is being uploaded, reducing the total duration of sequential transfers. ## Changes - **Pattern Implementation:** Refactored `uploadVideoChunks` into a producer-consumer model using a dedicated thread for pre-loading data from the input stream. - **Memory Management:** Introduced a `bufferPool` with a fixed size (2 chunks) to ensure predictable memory usage (~100MB) regardless of video size. - **Asynchronous Execution:** Utilized `CompletableFuture` and `LinkedBlockingQueue` for thread-safe coordination between chunk production and upload. - **Robustness:** Added explicit error propagation and resource cleanup (buffers, threads) using `AtomicBoolean` and `finally` blocks. - **Observability:** Enhanced logging with detailed markers for chunk fetching, queueing, and upload progress to facilitate monitoring and debugging. image --- .../synology/service/SynologyDTPService.java | 213 ++++++++++++++---- 1 file changed, 175 insertions(+), 38 deletions(-) diff --git a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java index bceb1fb28..56ae45dc9 100644 --- a/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java +++ b/extensions/data-transfer/portability-data-transfer-synology/src/main/java/org/datatransferproject/datatransfer/synology/service/SynologyDTPService.java @@ -29,6 +29,11 @@ import java.util.Date; import java.util.Map; import java.util.UUID; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import okhttp3.FormBody; import okhttp3.MediaType; import okhttp3.MultipartBody; @@ -271,60 +276,192 @@ public Map createVideo(VideoModel video, UUID jobId) private int uploadVideoChunks(VideoModel video, UUID jobId, InputStream inputStream) throws CopyExceptionWithFailureReason, IOException { final int CHUNK_SIZE = 50 * 1024 * 1024; // 50MB - int actualChunkCount = 0; + // Use a pool of 2 buffers to limit memory usage to ~100MB regardless of video + BlockingQueue bufferPool = new LinkedBlockingQueue<>(2); + bufferPool.add(new byte[CHUNK_SIZE]); + bufferPool.add(new byte[CHUNK_SIZE]); + + // Pass video chunks from producer to consumer using a blocking queue + // the producer will put a VideoChunk with null data to indicate the end of stream or an error. + BlockingQueue chunkQueue = new LinkedBlockingQueue<>(); + AtomicBoolean consumerAborted = new AtomicBoolean(false); + + // Start a producer thread to read the input stream. + CompletableFuture producerFuture = + CompletableFuture.runAsync( + () -> + produceVideoChunks( + video, + jobId, + inputStream, + bufferPool, + chunkQueue, + consumerAborted, + CHUNK_SIZE)); + + try { + return consumeAndUploadVideoChunks(video, jobId, bufferPool, chunkQueue); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IOException("Video upload interrupted", e); + } finally { + consumerAborted.set(true); + producerFuture.cancel(true); + } + } + + private void produceVideoChunks( + VideoModel video, + UUID jobId, + InputStream inputStream, + BlockingQueue bufferPool, + BlockingQueue chunkQueue, + AtomicBoolean consumerAborted, + int chunkSize) { + try { + int index = 0; + while (!consumerAborted.get()) { + byte[] chunkData = bufferPool.poll(500, TimeUnit.MILLISECONDS); + if (chunkData == null) { + continue; + } + + int bytesRead = ByteStreams.read(inputStream, chunkData, 0, chunkSize); + if (bytesRead == 0) { + bufferPool.add(chunkData); // Return unused buffer + chunkQueue.add(new VideoChunk(null, 0, -1)); + break; + } - // reuse the same byte array for each chunk to reduce memory usage, - // since we are uploading in sequential, there is no concurrency issue - byte[] chunkData = new byte[CHUNK_SIZE]; + final int fIndex = index; + monitor.info( + () -> + String.format( + "[SynologyImporter] fetched video chunk, video name: [%s], chunk index: [%d]," + + " chunk size: [%d].", + video.getName(), fIndex, bytesRead), + jobId); + + chunkQueue.add(new VideoChunk(chunkData, bytesRead, index++)); + + monitor.info( + () -> + String.format( + "[SynologyImporter] video chunk put into queue, video name: [%s], chunk index:" + + " [%d], chunk size: [%d].", + video.getName(), fIndex, bytesRead), + jobId); + } + } catch (Exception e) { + chunkQueue.add(new VideoChunk(e)); + } + } + + private int consumeAndUploadVideoChunks( + VideoModel video, + UUID jobId, + BlockingQueue bufferPool, + BlockingQueue chunkQueue) + throws CopyExceptionWithFailureReason, IOException, InterruptedException { + int actualChunkCount = 0; while (true) { - int bytesRead = ByteStreams.read(inputStream, chunkData, 0, CHUNK_SIZE); - if (bytesRead == 0) { - break; + VideoChunk chunk = chunkQueue.take(); + if (chunk.error != null) { + throw new IOException("Error reading video stream", chunk.error); + } + if (chunk.data == null) { + final int finalActualChunkCount = actualChunkCount; + monitor.info( + () -> + String.format( + "[SynologyImporter] finished reading video stream, total chunks: [%d].", + finalActualChunkCount), + jobId); + return actualChunkCount; } - RequestBody fileBody = - new RequestBody() { - @Override - public MediaType contentType() { - return MediaType.parse(video.getMimeType()); - } + monitor.info( + () -> + String.format( + "[SynologyImporter] get video chunk, video name: [%s], chunk index: [%d], chunk" + + " size: [%d].", + video.getName(), chunk.index, chunk.size), + jobId); - @Override - public long contentLength() { - return bytesRead; - } + try { + uploadVideoChunk(video, jobId, chunk); + } finally { + bufferPool.add(chunk.data); + } + actualChunkCount++; + } + } - @Override - public void writeTo(BufferedSink sink) throws IOException { - // write the actual bytes read for the last chunk, which may be smaller than - // CHUNK_SIZE - sink.write(chunkData, 0, bytesRead); - } - }; + private void uploadVideoChunk(VideoModel video, UUID jobId, VideoChunk chunk) + throws CopyExceptionWithFailureReason, IOException { + final int bytesRead = chunk.size; + final byte[] data = chunk.data; + + RequestBodyGenerator bodyGenerator = + () -> { + RequestBody fileBody = + new RequestBody() { + @Override + public MediaType contentType() { + return MediaType.parse(video.getMimeType()); + } - MultipartBody.Builder builder = - new MultipartBody.Builder() + @Override + public long contentLength() { + return bytesRead; + } + + @Override + public void writeTo(BufferedSink sink) throws IOException { + sink.write(data, 0, bytesRead); + } + }; + + return new MultipartBody.Builder() .setType(MultipartBody.FORM) .addFormDataPart("item_id", video.getDataId()) - .addFormDataPart("index", String.valueOf(actualChunkCount)) + .addFormDataPart("index", String.valueOf(chunk.index)) .addFormDataPart("job_id", jobId.toString()) .addFormDataPart("service", exportingService) .addFormDataPart("file_size", String.valueOf(bytesRead)) - .addFormDataPart("file", video.getName(), fileBody); + .addFormDataPart("file", video.getName(), fileBody) + .build(); + }; - RequestBody requestBody = builder.build(); + monitor.info( + () -> + String.format( + "[SynologyImporter] uploading video chunk, video name: [%s], chunk index: [%d]," + + " chunk size: [%d].", + video.getName(), chunk.index, bytesRead), + jobId); + sendPostRequest(c2Api.getChunkUploadItem(), bodyGenerator, jobId); + } - final String chunkInfo = - String.format( - "[SynologyImporter] uploading video chunk, video name: [%s], chunk index: [%d], chunk" - + " size: [%d].", - video.getName(), actualChunkCount, bytesRead); - monitor.info(() -> chunkInfo, jobId); - sendPostRequest(c2Api.getChunkUploadItem(), () -> requestBody, jobId); + private static class VideoChunk { + final byte[] data; + final int size; + final int index; + final Throwable error; + + VideoChunk(byte[] data, int size, int index) { + this.data = data; + this.size = size; + this.index = index; + this.error = null; + } - actualChunkCount++; + VideoChunk(Throwable error) { + this.data = null; + this.size = 0; + this.index = -1; + this.error = error; } - return actualChunkCount; } private Map completeVideoUpload(VideoModel video, UUID jobId, int totalChunks) From 452b6b1400c70055e5ca1d9e7abcd0b9085bf74a Mon Sep 17 00:00:00 2001 From: Yuktha Gaduputi <47736831+yukthagaduputi@users.noreply.github.com> Date: Tue, 19 May 2026 10:34:25 +0530 Subject: [PATCH 09/11] fix(transfer): Prevent JobPollingService thread crash on stale job claims (#1491) This PR fixes a critical process-terminating thread crash in the JobPollingService loop that occurs when a transfer worker attempts to claim a job that has already been claimed and has credentials stored by another worker instance. **Problem** In a distributed, highly concurrent deployment with multiple worker instances, a worker can poll a job ID from an eventually consistent database index that looks free (CREDS_AVAILABLE) but has actually already been claimed by a faster peer instance. When our worker performs a strongly consistent read (store.findJob) to retrieve the job details, it detects the instanceId mismatch (indicating the job belongs to another worker) and marks the job's state as CANCELED in memory. However, in the old implementation: 1. JobPollingService.tryToClaimJob failed to verify whether the retrieved existingJob was null or canceled, proceeding blindly to build the updatedJob to claim it. 2. Constructing this job object to transition to state CREDS_ENCRYPTION_KEY_GENERATED threw a validation exception (IllegalStateException) because credentials (encryptedAuthData) had already been set by the peer worker. 3. Because the PortabilityJob builder execution occurred outside the main try-catch block in tryToClaimJob, this exception was uncaught. 4. This uncaught thread failure propagated up, terminating the periodic JobPollingService thread and causing the entire container sandbox to crash. **Solution** Implemented two layers of safety in JobPollingService.java to resolve this: 1. Proactive Prevention (Early Abort): Added null and CANCELED state checks immediately after retrieving the job from the JobStore. If the job has been deleted or marked canceled (which happens on instanceId mismatch), the worker aborts the claim attempt early and returns false safely. 2. Reactive Safety (Validation Safety Net): Wrapped the PortabilityJob builder execution in a local try-catch block to safely handle any unexpected IllegalStateException validation errors during object construction, returning false (handled failure) instead of propagating and crashing the thread. These changes ensure that failing to claim a job (due to losing a race) is treated as a handled, temporary failure, allowing the worker to complete the current polling iteration normally and try again in the next cycle instead of crashing. --- .../transfer/JobPollingService.java | 48 ++++++++++++++----- 1 file changed, 36 insertions(+), 12 deletions(-) diff --git a/portability-transfer/src/main/java/org/datatransferproject/transfer/JobPollingService.java b/portability-transfer/src/main/java/org/datatransferproject/transfer/JobPollingService.java index ca06937e2..571ffcbca 100644 --- a/portability-transfer/src/main/java/org/datatransferproject/transfer/JobPollingService.java +++ b/portability-transfer/src/main/java/org/datatransferproject/transfer/JobPollingService.java @@ -152,6 +152,19 @@ private void pollForUnassignedJob() { private boolean tryToClaimJob(UUID jobId, WorkerKeyPair keyPair) { // Lookup the job so we can append to its existing properties. PortabilityJob existingJob = store.findJob(jobId); + if (existingJob == null) { + monitor.debug(() -> format("JobPollingService: tryToClaimJob: jobId: %s not found", jobId)); + return false; + } + if (existingJob.state() == PortabilityJob.State.CANCELED) { + monitor.debug( + () -> + format( + "JobPollingService: tryToClaimJob: jobId: %s is canceled (likely instance" + + " mismatch)", + jobId)); + return false; + } monitor.debug(() -> format("JobPollingService: tryToClaimJob: jobId: %s", existingJob)); // TODO: Consider moving this check earlier in the flow @@ -166,18 +179,29 @@ private boolean tryToClaimJob(UUID jobId, WorkerKeyPair keyPair) { } String serializedKey = publicKeySerializer.serialize(keyPair.getEncodedPublicKey()); - PortabilityJob updatedJob = - existingJob - .toBuilder() - .setAndValidateJobAuthorization( - existingJob - .jobAuthorization() - .toBuilder() - .setInstanceId(keyPair.getInstanceId()) - .setAuthPublicKey(serializedKey) - .setState(JobAuthorization.State.CREDS_ENCRYPTION_KEY_GENERATED) - .build()) - .build(); + PortabilityJob updatedJob; + try { + updatedJob = + existingJob.toBuilder() + .setAndValidateJobAuthorization( + existingJob + .jobAuthorization() + .toBuilder() + .setInstanceId(keyPair.getInstanceId()) + .setAuthPublicKey(serializedKey) + .setState(JobAuthorization.State.CREDS_ENCRYPTION_KEY_GENERATED) + .build()) + .build(); + } catch (IllegalStateException e) { + monitor.debug( + () -> + format( + "Could not build updated job object for %s. Validation failed. Error msg: %s", + jobId, e.getMessage()), + e); + return false; + } + // Attempt to 'claim' this job by validating it is still in state CREDS_AVAILABLE as we // update it to state CREDS_ENCRYPTION_KEY_GENERATED, along with our key. If another transfer // instance polled the same job, and already claimed it, it will have updated the job's state From 4f3ff8dea7b357c497f58f56c882f363ff7d8f95 Mon Sep 17 00:00:00 2001 From: ameya9 Date: Mon, 22 Jun 2026 15:29:59 +0100 Subject: [PATCH 10/11] fix(google-calendar): throw InvalidTokenException on invalid_grant Mirror DriveImporter's handling: when calendar/event insertion fails with a TokenResponseException whose error is "invalid_grant" (HTTP 400 from the oauth2 token endpoint), wrap and rethrow as InvalidTokenException so the transfer job recognizes the token as invalid instead of treating it as a generic failure. --- .../calendar/GoogleCalendarImporter.java | 39 ++++++++++++++----- 1 file changed, 29 insertions(+), 10 deletions(-) diff --git a/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java b/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java index 29d04d976..ebf161de6 100644 --- a/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java +++ b/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java @@ -1,6 +1,8 @@ package org.datatransferproject.datatransfer.google.calendar; import com.google.api.client.auth.oauth2.Credential; +import com.google.api.client.auth.oauth2.TokenErrorResponse; +import com.google.api.client.auth.oauth2.TokenResponseException; import com.google.api.client.util.DateTime; import com.google.api.services.calendar.Calendar; import com.google.api.services.calendar.model.Event; @@ -12,6 +14,7 @@ import org.datatransferproject.spi.transfer.idempotentexecutor.IdempotentImportExecutor; import org.datatransferproject.spi.transfer.provider.ImportResult; import org.datatransferproject.spi.transfer.provider.Importer; +import org.datatransferproject.spi.transfer.types.InvalidTokenException; import org.datatransferproject.types.common.models.calendar.CalendarAttendeeModel; import org.datatransferproject.types.common.models.calendar.CalendarContainerResource; import org.datatransferproject.types.common.models.calendar.CalendarEventModel; @@ -113,26 +116,42 @@ public ImportResult importItem(UUID jobId, @VisibleForTesting String importSingleCalendar(TokensAndUrlAuthData authData, CalendarModel calendarModel) - throws IOException { + throws IOException, InvalidTokenException { com.google.api.services.calendar.model.Calendar toInsert = convertToGoogleCalendar( calendarModel); - com.google.api.services.calendar.model.Calendar calendarResult = - getOrCreateCalendarInterface(authData).calendars().insert(toInsert).execute(); - return calendarResult.getId(); + try { + com.google.api.services.calendar.model.Calendar calendarResult = + getOrCreateCalendarInterface(authData).calendars().insert(toInsert).execute(); + return calendarResult.getId(); + } catch (TokenResponseException e) { + TokenErrorResponse details = e.getDetails(); + if (details != null && details.getError().equals("invalid_grant")) { + throw new InvalidTokenException("Unable to refresh token.", e); + } + throw e; + } } @VisibleForTesting String importSingleEvent(IdempotentImportExecutor idempotentImportExecutor, TokensAndUrlAuthData authData, CalendarEventModel eventModel) - throws IOException { + throws IOException, InvalidTokenException { Event event = convertToGoogleCalendarEvent(eventModel); String newCalendarId = idempotentImportExecutor.getCachedValue(eventModel.getCalendarId()); - return getOrCreateCalendarInterface(authData) - .events() - .insert(newCalendarId, event) - .execute() - .getId(); + try { + return getOrCreateCalendarInterface(authData) + .events() + .insert(newCalendarId, event) + .execute() + .getId(); + } catch (TokenResponseException e) { + TokenErrorResponse details = e.getDetails(); + if (details != null && details.getError().equals("invalid_grant")) { + throw new InvalidTokenException("Unable to refresh token.", e); + } + throw e; + } } private Calendar getOrCreateCalendarInterface(TokensAndUrlAuthData authData) { From db5b4abbe4f8eab13460b2b5ef8c4eccc631201e Mon Sep 17 00:00:00 2001 From: Ameya Shendre <52448570+ameya9@users.noreply.github.com> Date: Tue, 7 Jul 2026 15:03:38 +0100 Subject: [PATCH 11/11] Apply suggestions from code review Co-authored-by: Alex Kulikov <7394728+alexeyqu@users.noreply.github.com> --- .../datatransfer/google/calendar/GoogleCalendarImporter.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java b/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java index ebf161de6..d8084dc1a 100644 --- a/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java +++ b/extensions/data-transfer/portability-data-transfer-google/src/main/java/org/datatransferproject/datatransfer/google/calendar/GoogleCalendarImporter.java @@ -125,7 +125,7 @@ String importSingleCalendar(TokensAndUrlAuthData authData, CalendarModel calenda return calendarResult.getId(); } catch (TokenResponseException e) { TokenErrorResponse details = e.getDetails(); - if (details != null && details.getError().equals("invalid_grant")) { + if (details != null && "invalid_grant".equals(details.getError())) { throw new InvalidTokenException("Unable to refresh token.", e); } throw e; @@ -147,7 +147,7 @@ String importSingleEvent(IdempotentImportExecutor idempotentImportExecutor, .getId(); } catch (TokenResponseException e) { TokenErrorResponse details = e.getDetails(); - if (details != null && details.getError().equals("invalid_grant")) { + if (details != null && "invalid_grant".equals(details.getError())) { throw new InvalidTokenException("Unable to refresh token.", e); } throw e;