committed by
GitHub
21 changed files with 663 additions and 178 deletions
@ -0,0 +1,126 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.entitiy.dashboard; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.ResourceType; |
|||
import org.thingsboard.server.common.data.TbResource; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.dao.resource.ResourceService; |
|||
import org.thingsboard.server.queue.discovery.PartitionService; |
|||
import org.thingsboard.server.queue.util.AfterStartUp; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.entitiy.widgets.bundle.TbWidgetsBundleService; |
|||
import org.thingsboard.server.service.sync.GitSyncService; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository.FileType; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository.RepoFile; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
import java.util.List; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.stream.Stream; |
|||
|
|||
@Service |
|||
@TbCoreComponent |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class DashboardSyncService { |
|||
|
|||
private final GitSyncService gitSyncService; |
|||
private final ResourceService resourceService; |
|||
private final TbWidgetsBundleService tbWidgetsBundleService; |
|||
private final PartitionService partitionService; |
|||
|
|||
@Value("${transport.gateway.dashboard.sync.enabled:true}") |
|||
private boolean enabled; |
|||
@Value("${transport.gateway.dashboard.sync.repository_url:}") |
|||
private String repoUrl; |
|||
@Value("${transport.gateway.dashboard.sync.fetch_frequency:24}") |
|||
private int fetchFrequencyHours; |
|||
|
|||
private static final String REPO_KEY = "gateways-dashboard"; |
|||
private static final String GATEWAY_RESOURCE_ID_PARAM = "${GATEWAY_RESOURCE_ID}"; |
|||
private static final String GATEWAYS_DASHBOARD_KEY = "gateways_dashboard.json"; |
|||
|
|||
@AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) |
|||
public void init() throws Exception { |
|||
if (!enabled) { |
|||
return; |
|||
} |
|||
gitSyncService.registerSync(REPO_KEY, repoUrl, "main", TimeUnit.HOURS.toMillis(fetchFrequencyHours), this::update); |
|||
} |
|||
|
|||
private void update() { |
|||
if (!partitionService.isMyPartition(ServiceType.TB_CORE, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID)) { |
|||
return; |
|||
} |
|||
|
|||
RepoFile extensionResourceFile = listFiles("resources").get(0); |
|||
String data = getFileContent(extensionResourceFile.path()); |
|||
TbResource extensionResource = createOrUpdateResource(ResourceType.JS_MODULE, extensionResourceFile.name(), data.getBytes(StandardCharsets.UTF_8)); |
|||
String extensionResourceId = extensionResource.getUuidId().toString(); |
|||
|
|||
Stream<JsonNode> widgetsBundles = listFiles("widget_bundles").stream() |
|||
.map(widgetsBundleFile -> { |
|||
String widgetsBundleDescriptor = getFileContent(widgetsBundleFile.path()); |
|||
widgetsBundleDescriptor = widgetsBundleDescriptor.replace(GATEWAY_RESOURCE_ID_PARAM, extensionResourceId); |
|||
return JacksonUtil.toJsonNode(widgetsBundleDescriptor); |
|||
}); |
|||
Stream<JsonNode> widgetTypes = listFiles("widget_types").stream() |
|||
.map(widgetTypeFile -> { |
|||
String widgetTypeDetails = getFileContent(widgetTypeFile.path()); |
|||
widgetTypeDetails = widgetTypeDetails.replace(GATEWAY_RESOURCE_ID_PARAM, extensionResourceId); |
|||
return JacksonUtil.toJsonNode(widgetTypeDetails); |
|||
}); |
|||
tbWidgetsBundleService.updateWidgets(TenantId.SYS_TENANT_ID, widgetsBundles, widgetTypes); |
|||
|
|||
RepoFile dashboardFile = listFiles("dashboards").get(0); |
|||
String dashboardJson = getFileContent(dashboardFile.path()).replace(GATEWAY_RESOURCE_ID_PARAM, extensionResourceId); |
|||
createOrUpdateResource(ResourceType.DASHBOARD, GATEWAYS_DASHBOARD_KEY, dashboardJson.getBytes(StandardCharsets.UTF_8)); |
|||
|
|||
log.info("Gateways dashboard sync completed"); |
|||
} |
|||
|
|||
private TbResource createOrUpdateResource(ResourceType resourceType, String resourceKey, byte[] data) { |
|||
TbResource resource = resourceService.findResourceByTenantIdAndKey(TenantId.SYS_TENANT_ID, resourceType, resourceKey); |
|||
if (resource == null) { |
|||
resource = new TbResource(); |
|||
resource.setTenantId(TenantId.SYS_TENANT_ID); |
|||
resource.setResourceType(resourceType); |
|||
resource.setResourceKey(resourceKey); |
|||
resource.setFileName(resourceKey); |
|||
resource.setTitle(resourceKey); |
|||
} |
|||
resource.setData(data); |
|||
log.debug("{} resource {}", (resource.getId() == null ? "Creating" : "Updating"), resourceKey); |
|||
return resourceService.saveResource(resource); |
|||
} |
|||
|
|||
private List<RepoFile> listFiles(String path) { |
|||
return gitSyncService.listFiles(REPO_KEY, path, 1, FileType.FILE); |
|||
} |
|||
|
|||
private String getFileContent(String path) { |
|||
return gitSyncService.getFileContent(REPO_KEY, path); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,168 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.sync; |
|||
|
|||
import jakarta.annotation.PreDestroy; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|||
import org.thingsboard.server.common.data.sync.vc.RepositorySettings; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository.FileType; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository.RepoFile; |
|||
|
|||
import java.net.URI; |
|||
import java.nio.file.Path; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.ScheduledExecutorService; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@TbCoreComponent |
|||
@Service |
|||
@Slf4j |
|||
public class DefaultGitSyncService implements GitSyncService { |
|||
|
|||
@Value("${vc.git.repositories-folder:${java.io.tmpdir}/repositories}") |
|||
private String repositoriesFolder; |
|||
|
|||
private final ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("git-sync")); |
|||
private final Map<String, GitRepository> repositories = new ConcurrentHashMap<>(); |
|||
private final Map<String, Runnable> updateListeners = new ConcurrentHashMap<>(); |
|||
|
|||
@Override |
|||
public void registerSync(String key, String repoUri, String branch, long fetchFrequencyMs, Runnable onUpdate) { |
|||
RepositorySettings settings = new RepositorySettings(); |
|||
settings.setRepositoryUri(repoUri); |
|||
settings.setDefaultBranch(branch); |
|||
if (onUpdate != null) { |
|||
updateListeners.put(key, onUpdate); |
|||
} |
|||
|
|||
executor.execute(() -> { |
|||
initRepository(key, settings); |
|||
}); |
|||
|
|||
executor.scheduleWithFixedDelay(() -> { |
|||
GitRepository repository = repositories.get(key); |
|||
if (repository == null || !GitRepository.exists(repository.getDirectory())) { |
|||
initRepository(key, settings); |
|||
return; |
|||
} |
|||
|
|||
try { |
|||
log.debug("[{}] Fetching repository", key); |
|||
repository.fetch(); |
|||
onUpdate(key); |
|||
} catch (Throwable e) { |
|||
log.error("[{}] Failed to fetch repository", key, e); |
|||
} |
|||
}, fetchFrequencyMs, fetchFrequencyMs, TimeUnit.MILLISECONDS); |
|||
} |
|||
|
|||
@Override |
|||
public List<RepoFile> listFiles(String key, String path, int depth, FileType type) { |
|||
GitRepository repository = getRepository(key); |
|||
return repository.listFilesAtCommit(getBranchRef(repository), path, depth).stream() |
|||
.filter(file -> type == null || file.type() == type) |
|||
.toList(); |
|||
} |
|||
|
|||
|
|||
@Override |
|||
public String getFileContent(String key, String path) { |
|||
GitRepository repository = getRepository(key); |
|||
try { |
|||
return repository.getFileContentAtCommit(path, getBranchRef(repository)); |
|||
} catch (Exception e) { |
|||
log.warn("[{}] Failed to get file content for path {}: {}", key, path, e.getMessage()); |
|||
return "{}"; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String getGithubRawContentUrl(String key, String path) { |
|||
if (path == null) { |
|||
return ""; |
|||
} |
|||
RepositorySettings settings = getRepository(key).getSettings(); |
|||
return StringUtils.removeEnd(settings.getRepositoryUri(), ".git") + "/blob/" + settings.getDefaultBranch() + "/" + path + "?raw=true"; |
|||
} |
|||
|
|||
private GitRepository getRepository(String key) { |
|||
GitRepository repository = repositories.get(key); |
|||
if (repository != null) { |
|||
if (!GitRepository.exists(repository.getDirectory())) { |
|||
// reinitializing the repository because folder was deleted
|
|||
initRepository(key, repository.getSettings()); |
|||
} |
|||
} |
|||
|
|||
repository = repositories.get(key); |
|||
if (repository == null) { |
|||
throw new IllegalStateException(key + " repository is not initialized"); |
|||
} |
|||
return repository; |
|||
} |
|||
|
|||
private void initRepository(String key, RepositorySettings settings) { |
|||
try { |
|||
repositories.remove(key); |
|||
Path directory = getRepoDirectory(settings); |
|||
|
|||
GitRepository repository = GitRepository.openOrClone(directory, settings, true); |
|||
repositories.put(key, repository); |
|||
log.info("[{}] Initialized repository", key); |
|||
|
|||
onUpdate(key); |
|||
} catch (Throwable e) { |
|||
log.error("[{}] Failed to initialize repository with settings {}", key, settings, e); |
|||
} |
|||
} |
|||
|
|||
private void onUpdate(String key) { |
|||
Runnable listener = updateListeners.get(key); |
|||
if (listener != null) { |
|||
log.debug("[{}] Handling repository update", key); |
|||
try { |
|||
listener.run(); |
|||
} catch (Throwable e) { |
|||
log.error("[{}] Failed to handle repository update", key, e); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private Path getRepoDirectory(RepositorySettings settings) { |
|||
// using uri to define folder name in case repo url is changed
|
|||
String name = URI.create(settings.getRepositoryUri()).getPath().replaceAll("[^a-zA-Z]", ""); |
|||
return Path.of(repositoriesFolder, name); |
|||
} |
|||
|
|||
private String getBranchRef(GitRepository repository) { |
|||
return "refs/remotes/origin/" + repository.getSettings().getDefaultBranch(); |
|||
} |
|||
|
|||
@PreDestroy |
|||
private void preDestroy() { |
|||
executor.shutdownNow(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,33 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.sync; |
|||
|
|||
import org.thingsboard.server.service.sync.vc.GitRepository.FileType; |
|||
import org.thingsboard.server.service.sync.vc.GitRepository.RepoFile; |
|||
|
|||
import java.util.List; |
|||
|
|||
public interface GitSyncService { |
|||
|
|||
void registerSync(String key, String repoUri, String branch, long fetchFrequencyMs, Runnable onUpdate); |
|||
|
|||
List<RepoFile> listFiles(String key, String path, int depth, FileType type); |
|||
|
|||
String getFileContent(String key, String path); |
|||
|
|||
String getGithubRawContentUrl(String key, String path); |
|||
|
|||
} |
|||
@ -0,0 +1,55 @@ |
|||
/** |
|||
* Copyright © 2016-2024 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.entitiy.dashboard; |
|||
|
|||
import org.junit.Test; |
|||
import org.springframework.mock.web.MockHttpServletResponse; |
|||
import org.springframework.test.context.TestPropertySource; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.Dashboard; |
|||
import org.thingsboard.server.controller.AbstractControllerTest; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.awaitility.Awaitility.await; |
|||
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
|||
|
|||
@DaoSqlTest |
|||
@TestPropertySource(properties = { |
|||
"transport.gateway.dashboard.sync.enabled=true" |
|||
}) |
|||
public class DashboardSyncServiceTest extends AbstractControllerTest { |
|||
|
|||
@Test |
|||
public void testGatewaysDashboardSync() throws Exception { |
|||
loginTenantAdmin(); |
|||
await().atMost(60, TimeUnit.SECONDS).untilAsserted(() -> { |
|||
MockHttpServletResponse response = doGet("/api/resource/dashboard/system/gateways_dashboard.json") |
|||
.andExpect(status().isOk()) |
|||
.andReturn().getResponse(); |
|||
String dashboardJson = response.getContentAsString(); |
|||
String etag = response.getHeader("ETag"); |
|||
|
|||
Dashboard dashboard = JacksonUtil.fromString(dashboardJson, Dashboard.class); |
|||
assertThat(dashboard).isNotNull(); |
|||
assertThat(dashboard.getTitle()).containsIgnoringCase("gateway"); |
|||
assertThat(etag).isNotBlank(); |
|||
}); |
|||
} |
|||
|
|||
} |
|||
Loading…
Reference in new issue