From 29612985c7ccc532f39f36ce13aa1d3eacc8fd95 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 1 Sep 2022 18:45:02 +0300 Subject: [PATCH] Minor improvements to the File Cache --- .../controller/OtaPackageController.java | 8 +- .../constructor/OtaPackageMsgConstructor.java | 1 + .../ota/DefaultTbOtaPackageService.java | 12 +- ...leImp.java => DefaultTbMultipartFile.java} | 17 +- .../transport/DefaultTransportApiService.java | 9 +- .../src/main/resources/thingsboard.yml | 6 +- common/cache/pom.xml | 5 - .../cache/ota/files/BaseFileCacheService.java | 158 +++++++++++++----- .../cache/ota/files/FileCacheService.java | 7 +- .../server/cache/ota/files/OtaFileState.java | 47 ++++++ .../cache/ota/files/TemporaryFileCleaner.java | 122 -------------- .../ota/service/BaseFileCacheServiceTest.java | 13 +- .../server/dao/ota/TbMultipartFile.java | 15 ++ .../server/dao/ota/BaseOtaPackageService.java | 26 ++- 14 files changed, 238 insertions(+), 208 deletions(-) rename application/src/main/java/org/thingsboard/server/service/ota/{TbMultipartFileImp.java => DefaultTbMultipartFile.java} (56%) create mode 100644 common/cache/src/main/java/org/thingsboard/server/cache/ota/files/OtaFileState.java delete mode 100644 common/cache/src/main/java/org/thingsboard/server/cache/ota/files/TemporaryFileCleaner.java diff --git a/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java b/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java index 5e1d377b83..015b16c132 100644 --- a/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java +++ b/application/src/main/java/org/thingsboard/server/controller/OtaPackageController.java @@ -185,13 +185,17 @@ public class OtaPackageController extends BaseController { OtaPackageId otaPackageId = new OtaPackageId(toUUID(strOtaPackageId)); OtaPackageInfo info = checkOtaPackageInfoId(otaPackageId, Operation.READ); ChecksumAlgorithm checksumAlgorithm = ChecksumAlgorithm.valueOf(checksumAlgorithmStr.toUpperCase()); - return tbOtaPackageService.saveOtaPackageData(info, checksum, checksumAlgorithm, file, getCurrentUser()); + try { + return tbOtaPackageService.saveOtaPackageData(info, checksum, checksumAlgorithm, file, getCurrentUser()); + }catch (Exception e){ + throw e; + } } @ApiOperation(value = "Get OTA Package Infos (getOtaPackages)", notes = "Returns a page of OTA Package Info objects owned by tenant. " + PAGE_DATA_PARAMETERS + OTA_PACKAGE_INFO_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH, - produces = "application/json") + produces = APPLICATION_JSON_VALUE) @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/otaPackages", method = RequestMethod.GET) @ResponseBody diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/OtaPackageMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/OtaPackageMsgConstructor.java index 42b47d6305..34a37f62fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/OtaPackageMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/OtaPackageMsgConstructor.java @@ -68,6 +68,7 @@ public class OtaPackageMsgConstructor { } if (otaPackage.getData() != null) { try { + //TODO: Refactor to avoid OOM on events to Edge builder.setData(ByteString.copyFrom(otaPackage.getData().readAllBytes())); } catch (IOException e){ throw new RuntimeException(e); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/ota/DefaultTbOtaPackageService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/ota/DefaultTbOtaPackageService.java index ad58df004a..5e9f1907c3 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/ota/DefaultTbOtaPackageService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/ota/DefaultTbOtaPackageService.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 The Thingsboard 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 - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

+ * + * http://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. @@ -34,7 +34,7 @@ import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.ota.util.ChecksumUtil; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; -import org.thingsboard.server.service.ota.TbMultipartFileImp; +import org.thingsboard.server.service.ota.DefaultTbMultipartFile; import java.io.IOException; @@ -89,7 +89,7 @@ public class DefaultTbOtaPackageService extends AbstractTbEntityService implemen otaPackage.setContentType(file.getContentType()); otaPackage.setData(file.getInputStream()); otaPackage.setDataSize(file.getSize()); - OtaPackageInfo savedOtaPackage = otaPackageService.saveOtaPackage(otaPackage, new TbMultipartFileImp(file)); + OtaPackageInfo savedOtaPackage = otaPackageService.saveOtaPackage(otaPackage, new DefaultTbMultipartFile(file)); notificationEntityService.notifyCreateOrUpdateOrDelete(tenantId, null, savedOtaPackage.getId(), savedOtaPackage, user, ActionType.UPDATED, true, null); return savedOtaPackage; diff --git a/application/src/main/java/org/thingsboard/server/service/ota/TbMultipartFileImp.java b/application/src/main/java/org/thingsboard/server/service/ota/DefaultTbMultipartFile.java similarity index 56% rename from application/src/main/java/org/thingsboard/server/service/ota/TbMultipartFileImp.java rename to application/src/main/java/org/thingsboard/server/service/ota/DefaultTbMultipartFile.java index 79cc158bea..937e13d122 100644 --- a/application/src/main/java/org/thingsboard/server/service/ota/TbMultipartFileImp.java +++ b/application/src/main/java/org/thingsboard/server/service/ota/DefaultTbMultipartFile.java @@ -1,3 +1,18 @@ +/** + * Copyright © 2016-2022 The Thingsboard 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 + * + * http://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.thingsboard.server.service.ota; import lombok.RequiredArgsConstructor; @@ -10,7 +25,7 @@ import java.io.InputStream; import java.util.Optional; @RequiredArgsConstructor -public class TbMultipartFileImp implements TbMultipartFile { +public class DefaultTbMultipartFile implements TbMultipartFile { @NotNull private final MultipartFile file; diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index b3b617d380..6a14a53161 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 The Thingsboard 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 - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

+ * + * http://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. @@ -604,6 +604,7 @@ public class DefaultTransportApiService implements TransportApiService { if (!otaPackageDataCache.has(otaPackageId.toString())) { OtaPackage otaPackage = otaPackageService.findOtaPackageById(tenantId, otaPackageId); try { + //TODO: Do not put to Redis/InMem Cache and use File system on the Transport service instead. otaPackageDataCache.put(otaPackageId.toString(), otaPackage.getData().readAllBytes()); } catch (IOException e) { log.error("Failed to cache ota package with id {}",otaPackage.getId(), e); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 3b921abdf2..a0fd477105 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -518,8 +518,8 @@ spring.resources.chain: content: enabled: "true" -spring.servlet.multipart.max-file-size: "50MB" -spring.servlet.multipart.max-request-size: "50MB" +spring.servlet.multipart.max-file-size: "512MB" +spring.servlet.multipart.max-request-size: "512MB" spring.jpa.properties.hibernate.jdbc.lob.non_contextual_creation: "true" spring.jpa.properties.hibernate.order_by.default_null_ordering: "${SPRING_JPA_PROPERTIES_HIBERNATE_ORDER_BY_DEFAULT_NULL_ORDERING:last}" @@ -1150,4 +1150,4 @@ management: include: '${METRICS_ENDPOINTS_EXPOSE:info}' files: - temporary_files_directory: "${TEMPORARY_FILES_DIRECTORY: ${java.io.tmpdir}}" \ No newline at end of file + temporary_files_directory: "${TEMPORARY_FILES_DIRECTORY:}" \ No newline at end of file diff --git a/common/cache/pom.xml b/common/cache/pom.xml index 38e68d11f4..d3f2c838c8 100644 --- a/common/cache/pom.xml +++ b/common/cache/pom.xml @@ -44,11 +44,6 @@ org.springframework.boot spring-boot-autoconfigure - - javax.annotation - javax.annotation-api - 1.3.2 - org.springframework.data spring-data-redis diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/BaseFileCacheService.java b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/BaseFileCacheService.java index 2132bfa0ac..679e674e33 100644 --- a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/BaseFileCacheService.java +++ b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/BaseFileCacheService.java @@ -18,83 +18,159 @@ package org.thingsboard.server.cache.ota.files; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.commons.io.FileUtils; +import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.context.event.ApplicationReadyEvent; +import org.springframework.context.event.EventListener; +import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.OtaPackageId; +import javax.annotation.PostConstruct; import java.io.File; import java.io.IOException; import java.io.InputStream; -import java.util.Optional; -import java.util.UUID; +import java.net.URI; +import java.nio.channels.FileChannel; +import java.nio.channels.FileLock; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.nio.file.StandardOpenOption; +import java.util.Arrays; +import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.locks.Lock; +import java.util.function.Supplier; +import java.util.stream.Collectors; @Slf4j @RequiredArgsConstructor @Component public class BaseFileCacheService implements FileCacheService { - @Value("${files.temporary_files_directory}/ota/") - private String PATH; + + @Value("${files.temporary_files_directory:}") + private String tmpDir; + @Value("${java.io.tmpdir}") + private String defaultTmpDir; + private final static String FILE_NAME_TEMPLATE = "%s.tmp"; - private final TemporaryFileCleaner fileCleaner; - private final ConcurrentMap files = new ConcurrentHashMap<>(); + private final static long TEMPORARY_FILE_INACTIVITY_TIME = 900_000; + private final ConcurrentMap files = new ConcurrentHashMap<>(); + @PostConstruct + public void init() { + if (StringUtils.isEmpty(tmpDir)) { + tmpDir = defaultTmpDir; + } + } @Override - public File saveDataTemporaryFile(InputStream inputStream) { - File path = new File(PATH); + public File findOrLoad(OtaPackageId otaId, Supplier data) { + if (otaId == null || data == null) { + log.error("Received null variables: {}", otaId == null ? "otaPackageId" : "data"); + throw new RuntimeException("Input values can not be null"); + } + var state = files.computeIfAbsent(otaId, tmp -> + new OtaFileState(otaId, Paths.get(tmpDir, "ota", String.format(FILE_NAME_TEMPLATE, otaId.getId().toString())))); + Lock lock = state.getLock(); + lock.lock(); try { - File tempFile = File.createTempFile(UUID.randomUUID().toString(), ".tmp", path); - FileUtils.copyInputStreamToFile(inputStream, tempFile); - return tempFile; - } catch (IOException e) { - log.error("Failed to create temp file", e); - throw new RuntimeException("Failed to create temp file for input stream"); + if (!state.exists()) { + saveAsSystemFile(state.getFile(), data.get()); + } + state.updateLastActivityTime(); + return state.getFile(); + } finally { + lock.unlock(); } } - @Override - public Optional getOtaDataFile(OtaPackageId otaPackageId) { - String fileName = PATH + String.format(FILE_NAME_TEMPLATE,otaPackageId.getId().toString()); - if (exist(fileName)) { - fileCleaner.updateFileUsageStatus(otaPackageId); - return Optional.of(new File(fileName)); + @EventListener(ApplicationReadyEvent.class) + public void cleanDirectoryWithTemporaryFiles() { + createTempDirectoryIfNotExist(); + cleanDirectory(); + log.info("Directory {} with temporary ota files cleaned", tmpDir); + } + + private void createTempDirectoryIfNotExist() { + File directory = Paths.get(tmpDir, "ota").toFile(); + if (!directory.exists()) { + try { + FileUtils.forceMkdir(directory); + } catch (IOException e) { + log.error("Failed to create directory for temporary files ", e); + } } - return Optional.empty(); } - @Override - public File loadToFile(OtaPackageId otaPackageId, InputStream data) { - if (otaPackageId == null || data == null) { - log.error("Received null variables: {}", otaPackageId == null ? "otaPackageId" : "data"); - throw new RuntimeException("Input values can not be null"); + private void cleanDirectory() { + File directory = Paths.get(tmpDir, "ota").toFile(); + if (directory.isDirectory()) { + File[] files = directory.listFiles(); + if (files == null) return; + Arrays.stream(files).forEach( + file -> { + try { + FileUtils.delete(file); + } catch (Exception e) { + log.error("Failed to delete file {}", file.getName(), e); + } + } + ); + } + } + + @Scheduled(fixedDelay = 600_000) + private void deleteUnusedTemporaryFiles() { + long currentTime = System.currentTimeMillis(); + List toBeDeleted = files.values() + .stream() + .filter(entry -> isExpired(currentTime, entry)) + .collect(Collectors.toList()); + try { + toBeDeleted.forEach(state -> { + state.getLock().lock(); + try { + if (isExpired(currentTime, state)) { + deleteFile(state.getFile()); + } + files.remove(state.getOtaId()); + } finally { + state.getLock().unlock(); + } + }); + } catch (Exception e) { + log.error("Failed to delete unused files", e); } - files.computeIfAbsent(otaPackageId, ota -> processFileSaving(ota, data)); - String fileName = PATH + String.format(FILE_NAME_TEMPLATE,otaPackageId.getId().toString()); - return new File(fileName); + log.info("Deleted {} unused temporary files", toBeDeleted.size()); } + private boolean isExpired(long currentTime, OtaFileState entry) { + return currentTime - entry.getLastActivityTime() > TEMPORARY_FILE_INACTIVITY_TIME; + } - private Boolean processFileSaving(OtaPackageId otaPackageId, InputStream data) { - String fileName = PATH + String.format(FILE_NAME_TEMPLATE,otaPackageId.getId().toString()); - saveAsSystemFile(fileName, data); - fileCleaner.updateFileUsageStatus(otaPackageId); - return true; + private synchronized void deleteFile(File file) { + try (FileChannel channel = FileChannel.open(Path.of(URI.create(file.getPath())), StandardOpenOption.APPEND)) { + FileLock lock = channel.lock(); + if (file.exists()) { + FileUtils.delete(file); + log.info("System file {} was deleted", file.getName()); + } + lock.release(); + } catch (IOException e) { + log.error("Failed to delete file {}", file.getName(), e); + } } - private void saveAsSystemFile(String fileName, InputStream inputStream) { + private void saveAsSystemFile(File file, InputStream inputStream) { try { - File file = new File(fileName); FileUtils.copyInputStreamToFile(inputStream, file); } catch (IOException e) { - log.error("Failed to copy stream to system file {}", fileName, e); + log.error("Failed to copy stream to system file {}", file.getName(), e); throw new RuntimeException("Failed to save file"); } } - private boolean exist(String name) { - File file = new File(name); - return file.exists(); - } } \ No newline at end of file diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/FileCacheService.java b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/FileCacheService.java index cda01b2014..6d9505b1b7 100644 --- a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/FileCacheService.java +++ b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/FileCacheService.java @@ -21,9 +21,10 @@ import java.io.File; import java.io.FileNotFoundException; import java.io.InputStream; import java.util.Optional; +import java.util.function.Supplier; public interface FileCacheService { - File loadToFile(OtaPackageId otaPackageId, InputStream data); - File saveDataTemporaryFile(InputStream inputStream); - Optional getOtaDataFile(OtaPackageId otaPackageId) throws FileNotFoundException; + + File findOrLoad(OtaPackageId otaPackageId, Supplier data); + } diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/OtaFileState.java b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/OtaFileState.java new file mode 100644 index 0000000000..41e2231479 --- /dev/null +++ b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/OtaFileState.java @@ -0,0 +1,47 @@ +/** + * Copyright © 2016-2022 The Thingsboard 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 + * + * http://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.thingsboard.server.cache.ota.files; + +import lombok.Data; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.OtaPackageId; + +import java.io.File; +import java.nio.file.Path; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + +@Data +public class OtaFileState { + + private final Lock lock = new ReentrantLock(); + private final OtaPackageId otaId; + private final Path filePath; + + private long lastActivityTime; + + public void updateLastActivityTime() { + lastActivityTime = System.currentTimeMillis(); + } + + public boolean exists(){ + return filePath.toFile().exists(); + } + + public File getFile() { + return filePath.toFile(); + } +} diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/TemporaryFileCleaner.java b/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/TemporaryFileCleaner.java deleted file mode 100644 index 92f19189bf..0000000000 --- a/common/cache/src/main/java/org/thingsboard/server/cache/ota/files/TemporaryFileCleaner.java +++ /dev/null @@ -1,122 +0,0 @@ -/** - * Copyright © 2016-2022 The Thingsboard 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 - * - * http://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.thingsboard.server.cache.ota.files; - -import lombok.extern.slf4j.Slf4j; -import org.apache.commons.io.FileUtils; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.boot.context.event.ApplicationReadyEvent; -import org.springframework.context.event.EventListener; -import org.springframework.scheduling.annotation.Scheduled; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.id.OtaPackageId; - -import java.io.File; -import java.io.IOException; -import java.net.URI; -import java.nio.channels.FileChannel; -import java.nio.channels.FileLock; -import java.nio.file.Path; -import java.nio.file.StandardOpenOption; -import java.util.Arrays; -import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.stream.Collectors; - -@Slf4j -@Component -public class TemporaryFileCleaner { - @Value("${files.temporary_files_directory}/ota/") - private String PATH; - private final static String FILE_NAME_TEMPLATE = "%s%s.tmp"; - private final static long TEMPORARY_FILE_INACTIVITY_TIME = 900_000; - private final ConcurrentMap lastActivityTimes = new ConcurrentHashMap<>(); - - public void updateFileUsageStatus(OtaPackageId otaPackageId) { - lastActivityTimes.put(otaPackageId, System.currentTimeMillis()); - } - - @EventListener(ApplicationReadyEvent.class) - public void cleanDirectoryWithTemporaryFiles() { - createTempDirectoryIfNotExist(); - cleanDirectory(); - log.info("Directory {} with temporary ota files cleaned", PATH); - } - - private void createTempDirectoryIfNotExist() { - File directory = new File(PATH); - if (!directory.exists()) { - try{ - FileUtils.forceMkdir(directory); - } catch(IOException e){ - log.error("Failed to create directory for temporary files ", e); - } - } - } - - private void cleanDirectory() { - File directory = new File(PATH); - if (directory.isDirectory()) { - File[] files = directory.listFiles(); - if (files == null) return; - Arrays.stream(files).forEach( - file -> { - try { - FileUtils.delete(file); - } catch (Exception e) { - log.error("Failed to delete file {}", file.getName(), e); - } - } - ); - } - } - - @Scheduled(fixedDelay = 600_000) - private void deleteUnusedTemporaryFiles() { - long currentTime = System.currentTimeMillis(); - List toBeDeleted = lastActivityTimes.entrySet() - .stream() - .filter(entry -> currentTime - entry.getValue() > TEMPORARY_FILE_INACTIVITY_TIME) - .map(Map.Entry::getKey) - .collect(Collectors.toList()); - try { - toBeDeleted.forEach(otaId -> { - deleteFile(otaId.getId().toString()); - lastActivityTimes.remove(otaId); - }); - } catch (Exception e) { - log.error("Failed to delete unused files", e); - } - log.info("Deleted {} unused temporary files", toBeDeleted.size()); - } - - private synchronized void deleteFile(String otaId) { - String fileName = String.format(FILE_NAME_TEMPLATE, PATH, otaId); - File file = new File(fileName); - try (FileChannel channel = FileChannel.open(Path.of(URI.create(file.getPath())), StandardOpenOption.APPEND)) { - FileLock lock = channel.lock(); - if (file.exists()) { - FileUtils.delete(file); - log.info("System file {} was deleted", file.getName()); - } - lock.release(); - } catch (IOException e) { - log.error("Failed to delete file {}", file.getName(), e); - } - } -} diff --git a/common/cache/src/test/java/org/thingsboard/server/cache/ota/service/BaseFileCacheServiceTest.java b/common/cache/src/test/java/org/thingsboard/server/cache/ota/service/BaseFileCacheServiceTest.java index 455269773c..3086851a1c 100644 --- a/common/cache/src/test/java/org/thingsboard/server/cache/ota/service/BaseFileCacheServiceTest.java +++ b/common/cache/src/test/java/org/thingsboard/server/cache/ota/service/BaseFileCacheServiceTest.java @@ -18,7 +18,6 @@ package org.thingsboard.server.cache.ota.service; import lombok.SneakyThrows; import org.junit.jupiter.api.Test; import org.thingsboard.server.cache.ota.files.BaseFileCacheService; -import org.thingsboard.server.cache.ota.files.TemporaryFileCleaner; import org.thingsboard.server.common.data.id.OtaPackageId; import java.io.ByteArrayInputStream; @@ -39,17 +38,17 @@ class BaseFileCacheServiceTest { private static final int ONE_MEGA_BYTE = 1_000_000; private final static OtaPackageId OTA_PACKAGE_ID = new OtaPackageId(UUID.randomUUID()); private final static InputStream DATA = new ByteArrayInputStream(FILE_FILLING.getBytes()); - BaseFileCacheService baseFileCacheService = new BaseFileCacheService(new TemporaryFileCleaner()); + BaseFileCacheService baseFileCacheService = new BaseFileCacheService(); @Test void testDataSavingWithNullInputStream() { - assertThrows(RuntimeException.class, () -> baseFileCacheService.loadToFile(OTA_PACKAGE_ID, null)); + assertThrows(RuntimeException.class, () -> baseFileCacheService.findOrLoad(OTA_PACKAGE_ID, null)); } @Test void testDataSavingWithNullOtaPackageId() { - assertThrows(RuntimeException.class, () -> baseFileCacheService.loadToFile(null, DATA)); + assertThrows(RuntimeException.class, () -> baseFileCacheService.findOrLoad(null, () -> DATA)); } @Test @@ -57,8 +56,8 @@ class BaseFileCacheServiceTest { void testMultiSavingDataToFile() { File directory = new File(PATH); int beginning = Objects.requireNonNull(directory.list()).length; - Thread thread1 = new Thread(() -> baseFileCacheService.loadToFile(OTA_PACKAGE_ID, DATA)); - Thread thread2 = new Thread(() -> baseFileCacheService.loadToFile(OTA_PACKAGE_ID, DATA)); + Thread thread1 = new Thread(() -> baseFileCacheService.findOrLoad(OTA_PACKAGE_ID, () -> DATA)); + Thread thread2 = new Thread(() -> baseFileCacheService.findOrLoad(OTA_PACKAGE_ID, () -> DATA)); thread1.start(); thread2.start(); thread1.join(); @@ -72,7 +71,7 @@ class BaseFileCacheServiceTest { @SneakyThrows void testCorrectDataSavingToFile() { String sha256 = calculateChecksumSHA256(new ByteArrayInputStream(FILE_FILLING.getBytes())); - File file = baseFileCacheService.loadToFile(OTA_PACKAGE_ID, DATA); + File file = baseFileCacheService.findOrLoad(OTA_PACKAGE_ID, () -> DATA); assertEquals(sha256, calculateChecksumSHA256(new FileInputStream(file))); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/TbMultipartFile.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/TbMultipartFile.java index 1b92c6405f..eee18c489a 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/TbMultipartFile.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/ota/TbMultipartFile.java @@ -1,3 +1,18 @@ +/** + * Copyright © 2016-2022 The Thingsboard 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 + * + * http://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.thingsboard.server.dao.ota; import java.io.IOException; diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java index a538fc0359..a85ffe6996 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 The Thingsboard 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 - *

- * http://www.apache.org/licenses/LICENSE-2.0 - *

+ * + * http://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. @@ -213,15 +213,13 @@ public class BaseOtaPackageService extends AbstractCachedEntityService otaDataFile = baseFileCacheService.getOtaDataFile(otaPackageId); - if (otaDataFile.isPresent()) { - return otaDataFile.get(); - } - OtaPackage otaPackage = findOtaPackageById(tenantId, otaPackageId); - if (otaPackage == null) { - log.error("Can't find otaPackage to download file {}", otaPackageId); - throw new RuntimeException("No such OtaPackageId"); - } - return baseFileCacheService.loadToFile(otaPackageId, otaPackage.getData()); + return baseFileCacheService.findOrLoad(otaPackageId, () -> { + OtaPackage otaPackage = findOtaPackageById(tenantId, otaPackageId); + if (otaPackage == null) { + log.error("Can't find otaPackage to download file {}", otaPackageId); + throw new RuntimeException("No such OtaPackageId"); + } + return otaPackage.getData(); + }); } }