@ -1,43 +1,150 @@ |
|||||
# ThingsBoard |
 |
||||
[](https://builds.thingsboard.io/viewType.html?buildTypeId=ThingsBoard_Build&guest=1) |
|
||||
|
<div align="center"> |
||||
|
|
||||
|
# Open-source IoT platform for data collection, processing, visualization, and device management. |
||||
|
|
||||
|
</div> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
|
||||
|
💡 [Get started](https://thingsboard.io/docs/getting-started-guides/helloworld/) • 🌐 [Website](https://thingsboard.io/) • 📚 [Documentation](https://thingsboard.io/docs/) • 📔 [Blog](https://thingsboard.io/blog/) • ▶️ [Live demo](https://demo.thingsboard.io/signup) • 🔗 [LinkedIn](https://www.linkedin.com/company/thingsboard/posts/?feedView=all) |
||||
|
|
||||
|
</div> |
||||
|
|
||||
|
## 🚀 Installation options |
||||
|
|
||||
|
* Install ThingsBoard [On-premise](https://thingsboard.io/docs/user-guide/install/installation-options/?ceInstallType=onPremise) |
||||
|
* Try [ThingsBoard Cloud](https://thingsboard.io/installations/) |
||||
|
* or [Use our Live demo](https://demo.thingsboard.io/signup) |
||||
|
|
||||
|
## 💡 Getting started with ThingsBoard |
||||
|
|
||||
|
Check out our [Getting Started guide](https://thingsboard.io/docs/getting-started-guides/helloworld/) or [watch the video](https://www.youtube.com/watch?v=80L0ubQLXsc) to learn the basics of ThingsBoard and create your first dashboard! You will learn to: |
||||
|
|
||||
|
* Connect devices to ThingsBoard |
||||
|
* Push data from devices to ThingsBoard |
||||
|
* Build real-time dashboards |
||||
|
* Create a Customer and assign the dashboard with them. |
||||
|
* Define thresholds and trigger alarms |
||||
|
* Set up notifications via email, SMS, mobile apps, or integrate with third-party services. |
||||
|
|
||||
|
## ✨ Features |
||||
|
|
||||
|
<table> |
||||
|
<tr> |
||||
|
<td width="50%" valign="top"> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<img src="https://github.com/user-attachments/assets/255cca4f-b111-44e8-99ea-0af55f8e3681" alt="Provision and manage devices and assets" width="378" /> |
||||
|
<h3>Provision and manage <br> devices and assets</h3> |
||||
|
</div> |
||||
|
<div align="center"> |
||||
|
<p>Provision, monitor and control your IoT entities in secure way using rich server-side APIs. Define relations between your devices, assets, customers or any other entities.</p> |
||||
|
</div> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<a href="https://thingsboard.io/docs/user-guide/entities-and-relations/">Read more ➜</a> |
||||
|
</div> |
||||
|
<br> |
||||
|
</td> |
||||
|
<td width="50%" valign="top"> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<img src="https://github.com/user-attachments/assets/24b41d10-150a-42dd-ab1a-32ac9b5978c1" alt="Collect and visualize your data" width="378" /> |
||||
|
<h3>Collect and visualize <br> your data</h3> |
||||
|
</div> |
||||
|
<div align="center"> |
||||
|
<p>Collect and store telemetry data in scalable and fault-tolerant way. Visualize your data with built-in or custom widgets and flexible dashboards. Share dashboards with your customers.</p> |
||||
|
</div> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<a href="https://thingsboard.io/iot-data-visualization/">Read more ➜</a> |
||||
|
</div> |
||||
|
<br> |
||||
|
</td> |
||||
|
</tr> |
||||
|
<tr> |
||||
|
<td width="50%" valign="top"> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<img src="https://github.com/user-attachments/assets/6f2a6dd2-7b33-4d17-8b92-d1f995adda2c" alt="SCADA Dashboards" width="378" /> |
||||
|
<h3>SCADA Dashboards</h3> |
||||
|
</div> |
||||
|
<div align="center"> |
||||
|
<p>Monitor and control your industrial processes in real time with SCADA. Use SCADA symbols on dashboards to create and manage any workflow, offering full flexibility to design and oversee operations according to your requirements.</p> |
||||
|
</div> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<a href="https://thingsboard.io/use-cases/scada/">Read more ➜</a> |
||||
|
</div> |
||||
|
<br> |
||||
|
</td> |
||||
|
<td width="50%" valign="top"> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<img src="https://github.com/user-attachments/assets/c23dcc9b-aeba-40ef-9973-49b953fc1257" alt="Process and React" width="378" /> |
||||
|
<h3>Process and React</h3> |
||||
|
</div> |
||||
|
<div align="center"> |
||||
|
<p>Define data processing rule chains. Transform and normalize your device data. Raise alarms on incoming telemetry events, attribute updates, device inactivity and user actions.<br></p> |
||||
|
</div> |
||||
|
<br> |
||||
|
<br> |
||||
|
<div align="center"> |
||||
|
<a href="https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/">Read more ➜</a> |
||||
|
</div> |
||||
|
<br> |
||||
|
</td> |
||||
|
</tr> |
||||
|
</table> |
||||
|
|
||||
|
## ⚙️ Powerful IoT Rule Engine |
||||
|
|
||||
|
ThingsBoard allows you to create complex [Rule Chains](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) to process data from your devices and match your application specific use cases. |
||||
|
|
||||
|
[](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
||||
|
|
||||
|
<div align="center"> |
||||
|
|
||||
|
[**Read more about Rule Engine ➜**](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
||||
|
|
||||
|
</div> |
||||
|
|
||||
|
## 📦 Real-Time IoT Dashboards |
||||
|
|
||||
|
ThingsBoard is a scalable, user-friendly, and device-agnostic IoT platform that speeds up time-to-market with powerful built-in solution templates. It enables data collection and analysis from any devices, saving resources on routine tasks and letting you focus on your solution’s unique aspects. See more our Use Cases [here](https://thingsboard.io/iot-use-cases/). |
||||
|
|
||||
|
[**Smart energy**](https://thingsboard.io/use-cases/smart-energy/) |
||||
|
|
||||
|
[](https://thingsboard.io/use-cases/smart-energy/) |
||||
|
|
||||
|
[**SCADA swimming pool**](https://thingsboard.io/use-cases/scada/) |
||||
|
|
||||
|
[](https://thingsboard.io/use-cases/scada/) |
||||
|
|
||||
|
[**Fleet tracking**](https://thingsboard.io/use-cases/fleet-tracking/) |
||||
|
|
||||
|
[](https://thingsboard.io/use-cases/fleet-tracking/) |
||||
|
|
||||
|
[**Smart farming**](https://thingsboard.io/use-cases/smart-farming/) |
||||
|
|
||||
|
[](https://thingsboard.io/use-cases/smart-farming/) |
||||
|
|
||||
ThingsBoard is an open-source IoT platform for data collection, processing, visualization, and device management. |
[**Smart metering**](https://thingsboard.io/smart-metering/) |
||||
|
|
||||
<img src="./img/logo.png?raw=true" width="100" height="100"> |
|
||||
|
|
||||
|
|
||||
## Documentation |
|
||||
|
|
||||
ThingsBoard documentation is hosted on [thingsboard.io](https://thingsboard.io/docs). |
|
||||
|
|
||||
## IoT use cases |
|
||||
|
|
||||
[**Smart energy**](https://thingsboard.io/smart-energy/) |
|
||||
[](https://thingsboard.io/smart-energy/) |
|
||||
|
|
||||
[**SCADA Swimming pool**](https://thingsboard.io/use-cases/scada/) |
|
||||
[](https://thingsboard.io/use-cases/scada/) |
|
||||
|
|
||||
[**Fleet tracking**](https://thingsboard.io/fleet-tracking/) |
|
||||
[](https://thingsboard.io/fleet-tracking/) |
|
||||
|
|
||||
[**Smart farming**](https://thingsboard.io/smart-farming/) |
|
||||
[](https://thingsboard.io/smart-farming/) |
|
||||
|
|
||||
[**IoT Rule Engine**](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
[](https://thingsboard.io/smart-metering/) |
||||
[](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
|
||||
|
|
||||
[**Smart metering**](https://thingsboard.io/smart-metering/) |
<div align="center"> |
||||
[](https://thingsboard.io/smart-metering/) |
|
||||
|
|
||||
## Getting Started |
[**Check more of our use cases ➜**](https://thingsboard.io/iot-use-cases/) |
||||
|
|
||||
Collect and Visualize your IoT data in minutes by following this [guide](https://thingsboard.io/docs/getting-started-guides/helloworld/). |
</div> |
||||
|
|
||||
## Support |
## 🫶 Support |
||||
|
|
||||
- [Stackoverflow](http://stackoverflow.com/questions/tagged/thingsboard) |
To get support, please visit our [GitHub issues page](https://github.com/thingsboard/thingsboard/issues) |
||||
|
|
||||
## Licenses |
## 📄 Licenses |
||||
|
|
||||
This project is released under [Apache 2.0 License](./LICENSE). |
This project is released under [Apache 2.0 License](./LICENSE) |
||||
|
|||||
|
Before Width: | Height: | Size: 25 KiB After Width: | Height: | Size: 25 KiB |
|
Before Width: | Height: | Size: 18 KiB After Width: | Height: | Size: 18 KiB |
|
Before Width: | Height: | Size: 83 KiB After Width: | Height: | Size: 83 KiB |
|
Before Width: | Height: | Size: 11 KiB After Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 19 KiB After Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 16 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 111 KiB After Width: | Height: | Size: 111 KiB |
|
Before Width: | Height: | Size: 26 KiB After Width: | Height: | Size: 26 KiB |
|
Before Width: | Height: | Size: 25 KiB After Width: | Height: | Size: 25 KiB |
|
Before Width: | Height: | Size: 137 KiB After Width: | Height: | Size: 137 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 18 KiB |
|
Before Width: | Height: | Size: 31 KiB After Width: | Height: | Size: 32 KiB |
|
Before Width: | Height: | Size: 24 KiB After Width: | Height: | Size: 24 KiB |
|
Before Width: | Height: | Size: 16 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 15 KiB After Width: | Height: | Size: 15 KiB |
|
Before Width: | Height: | Size: 113 KiB After Width: | Height: | Size: 113 KiB |
|
Before Width: | Height: | Size: 84 KiB After Width: | Height: | Size: 84 KiB |
|
Before Width: | Height: | Size: 111 KiB After Width: | Height: | Size: 111 KiB |
|
Before Width: | Height: | Size: 118 KiB After Width: | Height: | Size: 118 KiB |
|
Before Width: | Height: | Size: 119 KiB After Width: | Height: | Size: 119 KiB |
|
Before Width: | Height: | Size: 112 KiB After Width: | Height: | Size: 112 KiB |
|
Before Width: | Height: | Size: 57 KiB After Width: | Height: | Size: 57 KiB |
|
Before Width: | Height: | Size: 54 KiB After Width: | Height: | Size: 54 KiB |
|
Before Width: | Height: | Size: 58 KiB After Width: | Height: | Size: 58 KiB |
|
Before Width: | Height: | Size: 72 KiB After Width: | Height: | Size: 72 KiB |
|
Before Width: | Height: | Size: 57 KiB After Width: | Height: | Size: 57 KiB |
|
Before Width: | Height: | Size: 54 KiB After Width: | Height: | Size: 54 KiB |
|
Before Width: | Height: | Size: 30 KiB After Width: | Height: | Size: 30 KiB |
|
Before Width: | Height: | Size: 42 KiB After Width: | Height: | Size: 42 KiB |
|
Before Width: | Height: | Size: 22 KiB After Width: | Height: | Size: 22 KiB |
|
Before Width: | Height: | Size: 22 KiB After Width: | Height: | Size: 22 KiB |
|
Before Width: | Height: | Size: 99 KiB After Width: | Height: | Size: 99 KiB |
|
Before Width: | Height: | Size: 46 KiB After Width: | Height: | Size: 46 KiB |
|
Before Width: | Height: | Size: 44 KiB After Width: | Height: | Size: 44 KiB |
|
Before Width: | Height: | Size: 46 KiB After Width: | Height: | Size: 46 KiB |
|
Before Width: | Height: | Size: 104 KiB After Width: | Height: | Size: 104 KiB |
|
Before Width: | Height: | Size: 117 KiB After Width: | Height: | Size: 117 KiB |
|
Before Width: | Height: | Size: 118 KiB After Width: | Height: | Size: 118 KiB |
|
Before Width: | Height: | Size: 120 KiB After Width: | Height: | Size: 120 KiB |
|
Before Width: | Height: | Size: 107 KiB After Width: | Height: | Size: 107 KiB |
|
Before Width: | Height: | Size: 119 KiB After Width: | Height: | Size: 119 KiB |
|
Before Width: | Height: | Size: 27 KiB After Width: | Height: | Size: 27 KiB |
|
Before Width: | Height: | Size: 22 KiB After Width: | Height: | Size: 22 KiB |
|
Before Width: | Height: | Size: 15 KiB After Width: | Height: | Size: 15 KiB |
|
Before Width: | Height: | Size: 100 KiB After Width: | Height: | Size: 100 KiB |
|
Before Width: | Height: | Size: 112 KiB After Width: | Height: | Size: 112 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 17 KiB |
|
Before Width: | Height: | Size: 21 KiB After Width: | Height: | Size: 22 KiB |
@ -0,0 +1,178 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.controller; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import dev.langchain4j.model.chat.request.ChatRequest; |
||||
|
import io.swagger.v3.oas.annotations.Parameter; |
||||
|
import io.swagger.v3.oas.annotations.media.Schema; |
||||
|
import jakarta.validation.Valid; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import org.springframework.security.access.prepost.PreAuthorize; |
||||
|
import org.springframework.validation.annotation.Validated; |
||||
|
import org.springframework.web.bind.annotation.DeleteMapping; |
||||
|
import org.springframework.web.bind.annotation.GetMapping; |
||||
|
import org.springframework.web.bind.annotation.PathVariable; |
||||
|
import org.springframework.web.bind.annotation.PostMapping; |
||||
|
import org.springframework.web.bind.annotation.RequestBody; |
||||
|
import org.springframework.web.bind.annotation.RequestMapping; |
||||
|
import org.springframework.web.bind.annotation.RequestParam; |
||||
|
import org.springframework.web.bind.annotation.RestController; |
||||
|
import org.springframework.web.context.request.async.DeferredResult; |
||||
|
import org.thingsboard.server.common.data.ai.AiModel; |
||||
|
import org.thingsboard.server.common.data.ai.dto.TbChatRequest; |
||||
|
import org.thingsboard.server.common.data.ai.dto.TbChatResponse; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.exception.ThingsboardException; |
||||
|
import org.thingsboard.server.common.data.id.AiModelId; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.config.annotations.ApiOperation; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.ai.AiChatModelService; |
||||
|
import org.thingsboard.server.service.security.permission.Operation; |
||||
|
import org.thingsboard.server.service.security.permission.Resource; |
||||
|
|
||||
|
import java.time.Duration; |
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static com.google.common.util.concurrent.MoreExecutors.directExecutor; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.AI_MODEL_TEXT_SEARCH_DESCRIPTION; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.PAGE_DATA_PARAMETERS; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.PAGE_NUMBER_DESCRIPTION; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.PAGE_SIZE_DESCRIPTION; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.SORT_ORDER_DESCRIPTION; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.SORT_PROPERTY_DESCRIPTION; |
||||
|
import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHORITY_PARAGRAPH; |
||||
|
|
||||
|
@Validated |
||||
|
@RestController |
||||
|
@TbCoreComponent |
||||
|
@RequiredArgsConstructor |
||||
|
@RequestMapping("/api/ai/model") |
||||
|
class AiModelController extends BaseController { |
||||
|
|
||||
|
private final AiChatModelService aiChatModelService; |
||||
|
|
||||
|
@ApiOperation( |
||||
|
value = "Create or update AI model (saveAiModel)", |
||||
|
notes = "Creates or updates an AI model record.\n\n" + |
||||
|
"• **Create:** Omit the `id` to create a new record. The platform assigns a UUID to the new record and returns it in the `id` field of the response.\n\n" + |
||||
|
"• **Update:** Include an existing `id` to modify that record. If no matching record exists, the API responds with **404 Not Found**.\n\n" + |
||||
|
"Tenant ID for the AI model will be taken from the authenticated user making the request, regardless of any value provided in the request body." + |
||||
|
TENANT_AUTHORITY_PARAGRAPH |
||||
|
) |
||||
|
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
||||
|
@PostMapping |
||||
|
public AiModel saveAiModel(@RequestBody @Valid AiModel model) throws ThingsboardException { |
||||
|
var user = getCurrentUser(); |
||||
|
model.setTenantId(user.getTenantId()); |
||||
|
checkEntity(model.getId(), model, Resource.AI_MODEL); |
||||
|
return tbAiModelService.save(model, user); |
||||
|
} |
||||
|
|
||||
|
@ApiOperation( |
||||
|
value = "Get AI model by ID (getAiModelById)", |
||||
|
notes = "Fetches an AI model record by its `id`." + |
||||
|
TENANT_AUTHORITY_PARAGRAPH |
||||
|
) |
||||
|
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
||||
|
@GetMapping("/{modelUuid}") |
||||
|
public AiModel getAiModelById( |
||||
|
@Parameter( |
||||
|
description = "ID of the AI model record", |
||||
|
required = true, |
||||
|
example = "de7900d4-30e2-11f0-9cd2-0242ac120002" |
||||
|
) |
||||
|
@PathVariable UUID modelUuid |
||||
|
) throws ThingsboardException { |
||||
|
return checkAiModelId(new AiModelId(modelUuid), Operation.READ); |
||||
|
} |
||||
|
|
||||
|
@ApiOperation( |
||||
|
value = "Get AI models (getAiModels)", |
||||
|
notes = "Returns a page of AI models. " + |
||||
|
PAGE_DATA_PARAMETERS + TENANT_AUTHORITY_PARAGRAPH |
||||
|
) |
||||
|
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
||||
|
@GetMapping |
||||
|
public PageData<AiModel> getAiModels( |
||||
|
@Parameter(description = PAGE_SIZE_DESCRIPTION, required = true) |
||||
|
@RequestParam int pageSize, |
||||
|
@Parameter(description = PAGE_NUMBER_DESCRIPTION, required = true) |
||||
|
@RequestParam int page, |
||||
|
@Parameter(description = AI_MODEL_TEXT_SEARCH_DESCRIPTION) |
||||
|
@RequestParam(required = false) String textSearch, |
||||
|
@Parameter(description = SORT_PROPERTY_DESCRIPTION, schema = @Schema(allowableValues = {"createdTime", "name", "provider", "modelId"})) |
||||
|
@RequestParam(required = false) String sortProperty, |
||||
|
@Parameter(description = SORT_ORDER_DESCRIPTION, schema = @Schema(allowableValues = {"ASC", "DESC"})) |
||||
|
@RequestParam(required = false) String sortOrder |
||||
|
) throws ThingsboardException { |
||||
|
var user = getCurrentUser(); |
||||
|
accessControlService.checkPermission(user, Resource.AI_MODEL, Operation.READ); |
||||
|
var pageLink = createPageLink(pageSize, page, textSearch, sortProperty, sortOrder); |
||||
|
return aiModelService.findAiModelsByTenantId(user.getTenantId(), pageLink); |
||||
|
} |
||||
|
|
||||
|
@ApiOperation( |
||||
|
value = "Delete AI model by ID (deleteAiModelById)", |
||||
|
notes = "Deletes the AI model record by its `id`. " + |
||||
|
"If a record with the specified `id` exists, the record is deleted and the endpoint returns `true`. " + |
||||
|
"If no such record exists, the endpoint returns `false`." + |
||||
|
TENANT_AUTHORITY_PARAGRAPH |
||||
|
) |
||||
|
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
||||
|
@DeleteMapping("/{modelUuid}") |
||||
|
public boolean deleteAiModelById( |
||||
|
@Parameter( |
||||
|
description = "ID of the AI model record", |
||||
|
required = true, |
||||
|
example = "de7900d4-30e2-11f0-9cd2-0242ac120002" |
||||
|
) |
||||
|
@PathVariable UUID modelUuid |
||||
|
) throws ThingsboardException { |
||||
|
var user = getCurrentUser(); |
||||
|
var modelId = new AiModelId(modelUuid); |
||||
|
accessControlService.checkPermission(user, Resource.AI_MODEL, Operation.DELETE); |
||||
|
Optional<AiModel> toDelete = aiModelService.findAiModelByTenantIdAndId(user.getTenantId(), modelId); |
||||
|
if (toDelete.isEmpty()) { |
||||
|
return false; |
||||
|
} |
||||
|
accessControlService.checkPermission(user, Resource.AI_MODEL, Operation.DELETE, modelId, toDelete.get()); |
||||
|
return tbAiModelService.delete(toDelete.get(), user); |
||||
|
} |
||||
|
|
||||
|
@ApiOperation( |
||||
|
value = "Send request to AI chat model (sendChatRequest)", |
||||
|
notes = "Submits a single prompt - made up of an optional system message and a required user message - to the specified AI chat model " + |
||||
|
"and returns either the generated answer or an error envelope." + |
||||
|
TENANT_AUTHORITY_PARAGRAPH |
||||
|
) |
||||
|
@PreAuthorize("hasAuthority('TENANT_ADMIN')") |
||||
|
@PostMapping("/chat") |
||||
|
public DeferredResult<TbChatResponse> sendChatRequest(@Valid @RequestBody TbChatRequest tbChatRequest) { |
||||
|
ChatRequest langChainChatRequest = tbChatRequest.toLangChainChatRequest(); |
||||
|
AiChatModelConfig<?> chatModelConfig = tbChatRequest.chatModelConfig(); |
||||
|
|
||||
|
ListenableFuture<TbChatResponse> future = aiChatModelService.sendChatRequestAsync(chatModelConfig, langChainChatRequest) |
||||
|
.transform(chatResponse -> (TbChatResponse) new TbChatResponse.Success(chatResponse.aiMessage().text()), directExecutor()) |
||||
|
.catching(Throwable.class, ex -> new TbChatResponse.Failure(ex.getMessage()), directExecutor()); |
||||
|
|
||||
|
Integer requestTimeoutSeconds = chatModelConfig.timeoutSeconds(); |
||||
|
return requestTimeoutSeconds != null ? wrapFuture(future, Duration.ofSeconds(requestTimeoutSeconds).toMillis()) : wrapFuture(future); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,20 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.ai; |
||||
|
|
||||
|
import org.thingsboard.rule.engine.api.RuleEngineAiChatModelService; |
||||
|
|
||||
|
public interface AiChatModelService extends RuleEngineAiChatModelService {} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.ai; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FluentFuture; |
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import dev.langchain4j.model.chat.request.ChatRequest; |
||||
|
import dev.langchain4j.model.chat.response.ChatResponse; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.Langchain4jChatModelConfigurer; |
||||
|
|
||||
|
@Service |
||||
|
@RequiredArgsConstructor |
||||
|
class AiChatModelServiceImpl implements AiChatModelService { |
||||
|
|
||||
|
private final Langchain4jChatModelConfigurer chatModelConfigurer; |
||||
|
private final AiRequestsExecutor aiRequestsExecutor; |
||||
|
|
||||
|
@Override |
||||
|
public <C extends AiChatModelConfig<C>> FluentFuture<ChatResponse> sendChatRequestAsync(AiChatModelConfig<C> chatModelConfig, ChatRequest chatRequest) { |
||||
|
ChatModel langChainChatModel = chatModelConfig.configure(chatModelConfigurer); |
||||
|
return aiRequestsExecutor.sendChatRequestAsync(langChainChatModel, chatRequest); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,27 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.ai; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FluentFuture; |
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import dev.langchain4j.model.chat.request.ChatRequest; |
||||
|
import dev.langchain4j.model.chat.response.ChatResponse; |
||||
|
|
||||
|
public interface AiRequestsExecutor { |
||||
|
|
||||
|
FluentFuture<ChatResponse> sendChatRequestAsync(ChatModel chatModel, ChatRequest chatRequest); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,86 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.ai; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FluentFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import dev.langchain4j.model.chat.request.ChatRequest; |
||||
|
import dev.langchain4j.model.chat.response.ChatResponse; |
||||
|
import jakarta.annotation.PostConstruct; |
||||
|
import jakarta.annotation.PreDestroy; |
||||
|
import jakarta.validation.constraints.Min; |
||||
|
import jakarta.validation.constraints.NotBlank; |
||||
|
import lombok.Data; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import org.springframework.boot.context.properties.ConfigurationProperties; |
||||
|
import org.springframework.context.annotation.Configuration; |
||||
|
import org.springframework.context.annotation.Lazy; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.springframework.validation.annotation.Validated; |
||||
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
||||
|
|
||||
|
import java.time.Duration; |
||||
|
import java.util.concurrent.Executors; |
||||
|
|
||||
|
@Lazy |
||||
|
@Component |
||||
|
@RequiredArgsConstructor |
||||
|
class DefaultAiRequestsExecutor implements AiRequestsExecutor { |
||||
|
|
||||
|
private final AiRequestsExecutorProperties properties; |
||||
|
|
||||
|
@Data |
||||
|
@Validated |
||||
|
@Configuration |
||||
|
@ConfigurationProperties(prefix = "actors.rule.ai-requests-thread-pool") |
||||
|
private static class AiRequestsExecutorProperties { |
||||
|
|
||||
|
@NotBlank(message = "Pool name must be not blank") |
||||
|
private String poolName = "ai-requests"; |
||||
|
|
||||
|
@Min(value = 1, message = "Pool size must be at least 1") |
||||
|
private int poolSize = 50; |
||||
|
|
||||
|
@Min(value = 1, message = "Termination timeout must be at least 1 second") |
||||
|
private int terminationTimeoutSeconds = 60; |
||||
|
|
||||
|
} |
||||
|
|
||||
|
private ListeningExecutorService executorService; |
||||
|
|
||||
|
@PostConstruct |
||||
|
private void init() { |
||||
|
executorService = MoreExecutors.listeningDecorator( |
||||
|
Executors.newFixedThreadPool(properties.getPoolSize(), ThingsBoardThreadFactory.forName(properties.getPoolName())) |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public FluentFuture<ChatResponse> sendChatRequestAsync(ChatModel chatModel, ChatRequest chatRequest) { |
||||
|
return FluentFuture.from(executorService.submit(() -> chatModel.chat(chatRequest))); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
private void destroy() { |
||||
|
if (executorService != null) { |
||||
|
MoreExecutors.shutdownAndAwaitTermination(executorService, Duration.ofSeconds(properties.getTerminationTimeoutSeconds())); |
||||
|
executorService = null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,269 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.ai; |
||||
|
|
||||
|
import com.google.api.gax.core.FixedCredentialsProvider; |
||||
|
import com.google.api.gax.retrying.RetrySettings; |
||||
|
import com.google.auth.oauth2.ServiceAccountCredentials; |
||||
|
import com.google.cloud.vertexai.Transport; |
||||
|
import com.google.cloud.vertexai.VertexAI; |
||||
|
import com.google.cloud.vertexai.api.GenerationConfig; |
||||
|
import com.google.cloud.vertexai.api.PredictionServiceClient; |
||||
|
import com.google.cloud.vertexai.api.PredictionServiceSettings; |
||||
|
import com.google.cloud.vertexai.generativeai.GenerativeModel; |
||||
|
import dev.langchain4j.model.anthropic.AnthropicChatModel; |
||||
|
import dev.langchain4j.model.azure.AzureOpenAiChatModel; |
||||
|
import dev.langchain4j.model.bedrock.BedrockChatModel; |
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import dev.langchain4j.model.chat.request.ChatRequestParameters; |
||||
|
import dev.langchain4j.model.github.GitHubModelsChatModel; |
||||
|
import dev.langchain4j.model.googleai.GoogleAiGeminiChatModel; |
||||
|
import dev.langchain4j.model.mistralai.MistralAiChatModel; |
||||
|
import dev.langchain4j.model.openai.OpenAiChatModel; |
||||
|
import dev.langchain4j.model.vertexai.gemini.VertexAiGeminiChatModel; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AmazonBedrockChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AnthropicChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.AzureOpenAiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.GitHubModelsChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.GoogleAiGeminiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.GoogleVertexAiGeminiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.Langchain4jChatModelConfigurer; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.MistralAiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.model.chat.OpenAiChatModelConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.AmazonBedrockProviderConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.AzureOpenAiProviderConfig; |
||||
|
import org.thingsboard.server.common.data.ai.provider.GoogleVertexAiGeminiProviderConfig; |
||||
|
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; |
||||
|
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; |
||||
|
import software.amazon.awssdk.regions.Region; |
||||
|
import software.amazon.awssdk.services.bedrockruntime.BedrockRuntimeClient; |
||||
|
|
||||
|
import java.io.ByteArrayInputStream; |
||||
|
import java.io.IOException; |
||||
|
import java.time.Duration; |
||||
|
|
||||
|
@Component |
||||
|
class Langchain4jChatModelConfigurerImpl implements Langchain4jChatModelConfigurer { |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(OpenAiChatModelConfig chatModelConfig) { |
||||
|
return OpenAiChatModel.builder() |
||||
|
.apiKey(chatModelConfig.providerConfig().apiKey()) |
||||
|
.modelName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.frequencyPenalty(chatModelConfig.frequencyPenalty()) |
||||
|
.presencePenalty(chatModelConfig.presencePenalty()) |
||||
|
.maxTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(AzureOpenAiChatModelConfig chatModelConfig) { |
||||
|
AzureOpenAiProviderConfig providerConfig = chatModelConfig.providerConfig(); |
||||
|
return AzureOpenAiChatModel.builder() |
||||
|
.endpoint(providerConfig.endpoint()) |
||||
|
.serviceVersion(providerConfig.serviceVersion()) |
||||
|
.apiKey(providerConfig.apiKey()) |
||||
|
.deploymentName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.frequencyPenalty(chatModelConfig.frequencyPenalty()) |
||||
|
.presencePenalty(chatModelConfig.presencePenalty()) |
||||
|
.maxTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(GoogleAiGeminiChatModelConfig chatModelConfig) { |
||||
|
return GoogleAiGeminiChatModel.builder() |
||||
|
.apiKey(chatModelConfig.providerConfig().apiKey()) |
||||
|
.modelName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.topK(chatModelConfig.topK()) |
||||
|
.frequencyPenalty(chatModelConfig.frequencyPenalty()) |
||||
|
.presencePenalty(chatModelConfig.presencePenalty()) |
||||
|
.maxOutputTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(GoogleVertexAiGeminiChatModelConfig chatModelConfig) { |
||||
|
GoogleVertexAiGeminiProviderConfig providerConfig = chatModelConfig.providerConfig(); |
||||
|
|
||||
|
// construct service account credentials using service account key JSON
|
||||
|
ServiceAccountCredentials serviceAccountCredentials; |
||||
|
try { |
||||
|
serviceAccountCredentials = ServiceAccountCredentials.fromStream(new ByteArrayInputStream(providerConfig.serviceAccountKey().getBytes())); |
||||
|
} catch (IOException e) { |
||||
|
throw new RuntimeException("Failed to parse service account key JSON", e); |
||||
|
} |
||||
|
|
||||
|
PredictionServiceSettings predictionServiceClientSettings; |
||||
|
try { |
||||
|
// create prediction service settings for REST transport with service account key credentials
|
||||
|
PredictionServiceSettings.Builder settingsBuilder = PredictionServiceSettings.newHttpJsonBuilder() |
||||
|
.setCredentialsProvider(FixedCredentialsProvider.create(serviceAccountCredentials)); |
||||
|
|
||||
|
// get the retry settings that control request timeout for generateContent RPC
|
||||
|
RetrySettings.Builder retrySettings = settingsBuilder |
||||
|
.generateContentSettings() |
||||
|
.getRetrySettings() |
||||
|
.toBuilder(); |
||||
|
|
||||
|
// set request timeout from model config
|
||||
|
if (chatModelConfig.timeoutSeconds() != null) { |
||||
|
retrySettings.setTotalTimeout(org.threeten.bp.Duration.ofSeconds(chatModelConfig.timeoutSeconds())); |
||||
|
} |
||||
|
|
||||
|
// set updated retry settings
|
||||
|
settingsBuilder.generateContentSettings().setRetrySettings(retrySettings.build()); |
||||
|
|
||||
|
// build the client settings
|
||||
|
predictionServiceClientSettings = settingsBuilder.build(); |
||||
|
} catch (IOException e) { |
||||
|
throw new RuntimeException("Failed to create prediction service client settings", e); |
||||
|
} |
||||
|
|
||||
|
// construct Vertex AI instance
|
||||
|
var vertexAI = new VertexAI.Builder() |
||||
|
.setProjectId(providerConfig.projectId()) |
||||
|
.setLocation(providerConfig.location()) |
||||
|
.setPredictionClientSupplier(() -> createPredictionServiceClient(predictionServiceClientSettings)) |
||||
|
.setTransport(Transport.REST) // GRPC also possible, but likely does not work with service account keys
|
||||
|
.build(); |
||||
|
|
||||
|
// map model config to generation config
|
||||
|
var generationConfigBuilder = GenerationConfig.newBuilder(); |
||||
|
if (chatModelConfig.temperature() != null) { |
||||
|
generationConfigBuilder.setTemperature(chatModelConfig.temperature().floatValue()); |
||||
|
} |
||||
|
if (chatModelConfig.topP() != null) { |
||||
|
generationConfigBuilder.setTopP(chatModelConfig.topP().floatValue()); |
||||
|
} |
||||
|
if (chatModelConfig.topK() != null) { |
||||
|
generationConfigBuilder.setTopK(chatModelConfig.topK()); |
||||
|
} |
||||
|
if (chatModelConfig.frequencyPenalty() != null) { |
||||
|
generationConfigBuilder.setFrequencyPenalty(chatModelConfig.frequencyPenalty().floatValue()); |
||||
|
} |
||||
|
if (chatModelConfig.frequencyPenalty() != null) { |
||||
|
generationConfigBuilder.setPresencePenalty(chatModelConfig.frequencyPenalty().floatValue()); |
||||
|
} |
||||
|
if (chatModelConfig.maxOutputTokens() != null) { |
||||
|
generationConfigBuilder.setMaxOutputTokens(chatModelConfig.maxOutputTokens()); |
||||
|
} |
||||
|
var generationConfig = generationConfigBuilder.build(); |
||||
|
|
||||
|
// construct generative model instance
|
||||
|
var generativeModel = new GenerativeModel(chatModelConfig.modelId(), vertexAI).withGenerationConfig(generationConfig); |
||||
|
|
||||
|
return new VertexAiGeminiChatModel(generativeModel, generationConfig, chatModelConfig.maxRetries()); |
||||
|
} |
||||
|
|
||||
|
private static PredictionServiceClient createPredictionServiceClient(PredictionServiceSettings settings) { |
||||
|
try { |
||||
|
return PredictionServiceClient.create(settings); |
||||
|
} catch (IOException e) { |
||||
|
throw new RuntimeException("Failed to create prediction service client", e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(MistralAiChatModelConfig chatModelConfig) { |
||||
|
return MistralAiChatModel.builder() |
||||
|
.apiKey(chatModelConfig.providerConfig().apiKey()) |
||||
|
.modelName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.frequencyPenalty(chatModelConfig.frequencyPenalty()) |
||||
|
.presencePenalty(chatModelConfig.presencePenalty()) |
||||
|
.maxTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(AnthropicChatModelConfig chatModelConfig) { |
||||
|
return AnthropicChatModel.builder() |
||||
|
.apiKey(chatModelConfig.providerConfig().apiKey()) |
||||
|
.modelName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.topK(chatModelConfig.topK()) |
||||
|
.maxTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(AmazonBedrockChatModelConfig chatModelConfig) { |
||||
|
AmazonBedrockProviderConfig providerConfig = chatModelConfig.providerConfig(); |
||||
|
|
||||
|
var credentialsProvider = StaticCredentialsProvider.create( |
||||
|
AwsBasicCredentials.create(providerConfig.accessKeyId(), providerConfig.secretAccessKey()) |
||||
|
); |
||||
|
|
||||
|
var bedrockClient = BedrockRuntimeClient.builder() |
||||
|
.region(Region.of(providerConfig.region())) |
||||
|
.credentialsProvider(credentialsProvider) |
||||
|
.build(); |
||||
|
|
||||
|
var defaultChatRequestParams = ChatRequestParameters.builder() |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.maxOutputTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.build(); |
||||
|
|
||||
|
return BedrockChatModel.builder() |
||||
|
.client(bedrockClient) |
||||
|
.modelId(chatModelConfig.modelId()) |
||||
|
.defaultRequestParameters(defaultChatRequestParams) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configureChatModel(GitHubModelsChatModelConfig chatModelConfig) { |
||||
|
return GitHubModelsChatModel.builder() |
||||
|
.gitHubToken(chatModelConfig.providerConfig().personalAccessToken()) |
||||
|
.modelName(chatModelConfig.modelId()) |
||||
|
.temperature(chatModelConfig.temperature()) |
||||
|
.topP(chatModelConfig.topP()) |
||||
|
.frequencyPenalty(chatModelConfig.frequencyPenalty()) |
||||
|
.presencePenalty(chatModelConfig.presencePenalty()) |
||||
|
.maxTokens(chatModelConfig.maxOutputTokens()) |
||||
|
.timeout(toDuration(chatModelConfig.timeoutSeconds())) |
||||
|
.maxRetries(chatModelConfig.maxRetries()) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
private static Duration toDuration(Integer timeoutSeconds) { |
||||
|
return timeoutSeconds != null ? Duration.ofSeconds(timeoutSeconds) : null; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,257 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.cf; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import jakarta.annotation.PostConstruct; |
||||
|
import jakarta.annotation.PreDestroy; |
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.kv.Aggregation; |
||||
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
||||
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
||||
|
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
||||
|
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
||||
|
import org.thingsboard.server.dao.attributes.AttributesService; |
||||
|
import org.thingsboard.server.dao.relation.RelationService; |
||||
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
||||
|
import org.thingsboard.server.dao.usagerecord.ApiLimitService; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; |
||||
|
|
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.Optional; |
||||
|
import java.util.Set; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; |
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
||||
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
||||
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; |
||||
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
||||
|
|
||||
|
@Data |
||||
|
@Slf4j |
||||
|
public abstract class AbstractCalculatedFieldProcessingService { |
||||
|
|
||||
|
protected final AttributesService attributesService; |
||||
|
protected final TimeseriesService timeseriesService; |
||||
|
protected final ApiLimitService apiLimitService; |
||||
|
protected final RelationService relationService; |
||||
|
|
||||
|
protected ListeningExecutorService calculatedFieldCallbackExecutor; |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( |
||||
|
Math.max(4, Runtime.getRuntime().availableProcessors()), getExecutorNamePrefix())); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void stop() { |
||||
|
if (calculatedFieldCallbackExecutor != null) { |
||||
|
calculatedFieldCallbackExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected abstract String getExecutorNamePrefix(); |
||||
|
|
||||
|
public ListenableFuture<CalculatedFieldState> fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) { |
||||
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCalculatedField().getType()) { |
||||
|
case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false); |
||||
|
case SIMPLE, SCRIPT -> { |
||||
|
Map<String, ListenableFuture<ArgumentEntry>> futures = new HashMap<>(); |
||||
|
for (var entry : ctx.getArguments().entrySet()) { |
||||
|
var argEntityId = resolveEntityId(entityId, entry.getValue()); |
||||
|
var argValueFuture = fetchArgumentValue(ctx.getTenantId(), argEntityId, entry.getValue(), System.currentTimeMillis()); |
||||
|
futures.put(entry.getKey(), argValueFuture); |
||||
|
} |
||||
|
yield futures; |
||||
|
} |
||||
|
}; |
||||
|
return Futures.whenAllComplete(argFutures.values()).call(() -> { |
||||
|
var result = createStateByType(ctx); |
||||
|
result.updateState(ctx, resolveArgumentFutures(argFutures)); |
||||
|
return result; |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected EntityId resolveEntityId(EntityId entityId, Argument argument) { |
||||
|
return argument.getRefEntityId() != null ? argument.getRefEntityId() : entityId; |
||||
|
} |
||||
|
|
||||
|
protected Map<String, ArgumentEntry> resolveArgumentFutures(Map<String, ListenableFuture<ArgumentEntry>> argFutures) { |
||||
|
return argFutures.entrySet().stream() |
||||
|
.collect(Collectors.toMap( |
||||
|
Map.Entry::getKey, // Keep the key as is
|
||||
|
entry -> { |
||||
|
try { |
||||
|
return entry.getValue().get(); |
||||
|
} catch (ExecutionException e) { |
||||
|
Throwable cause = e.getCause(); |
||||
|
throw new RuntimeException("Failed to fetch " + entry.getKey() + ": " + cause.getMessage(), cause); |
||||
|
} catch (InterruptedException e) { |
||||
|
throw new RuntimeException("Failed to fetch" + entry.getKey(), e); |
||||
|
} |
||||
|
} |
||||
|
)); |
||||
|
} |
||||
|
|
||||
|
protected Map<String, ListenableFuture<ArgumentEntry>> fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, boolean dynamicArgumentsOnly) { |
||||
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); |
||||
|
Set<Map.Entry<String, Argument>> entries = ctx.getArguments().entrySet(); |
||||
|
if (dynamicArgumentsOnly) { |
||||
|
entries = entries.stream() |
||||
|
.filter(entry -> entry.getValue().hasDynamicSource()) |
||||
|
.collect(Collectors.toSet()); |
||||
|
} |
||||
|
for (var entry : entries) { |
||||
|
switch (entry.getKey()) { |
||||
|
case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> |
||||
|
argFutures.put(entry.getKey(), fetchArgumentValue(ctx.getTenantId(), entityId, entry.getValue(), System.currentTimeMillis())); |
||||
|
default -> { |
||||
|
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry); |
||||
|
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds -> |
||||
|
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), MoreExecutors.directExecutor())); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return argFutures; |
||||
|
} |
||||
|
|
||||
|
private ListenableFuture<List<EntityId>> resolveGeofencingEntityIds(TenantId tenantId, EntityId entityId, Map.Entry<String, Argument> entry) { |
||||
|
Argument value = entry.getValue(); |
||||
|
if (value.getRefEntityId() != null) { |
||||
|
return Futures.immediateFuture(List.of(value.getRefEntityId())); |
||||
|
} |
||||
|
if (!value.hasDynamicSource()) { |
||||
|
return Futures.immediateFuture(List.of(entityId)); |
||||
|
} |
||||
|
var refDynamicSourceConfiguration = value.getRefDynamicSourceConfiguration(); |
||||
|
return switch (refDynamicSourceConfiguration.getType()) { |
||||
|
case RELATION_QUERY -> { |
||||
|
var configuration = (RelationQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; |
||||
|
if (configuration.isSimpleRelation()) { |
||||
|
yield switch (configuration.getDirection()) { |
||||
|
case FROM -> |
||||
|
Futures.transform(relationService.findByFromAndTypeAsync(tenantId, entityId, configuration.getRelationType(), RelationTypeGroup.COMMON), |
||||
|
configuration::resolveEntityIds, calculatedFieldCallbackExecutor); |
||||
|
case TO -> |
||||
|
Futures.transform(relationService.findByToAndTypeAsync(tenantId, entityId, configuration.getRelationType(), RelationTypeGroup.COMMON), |
||||
|
configuration::resolveEntityIds, calculatedFieldCallbackExecutor); |
||||
|
}; |
||||
|
} |
||||
|
yield Futures.transform(relationService.findByQuery(tenantId, configuration.toEntityRelationsQuery(entityId)), |
||||
|
configuration::resolveEntityIds, calculatedFieldCallbackExecutor); |
||||
|
} |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, Argument argument) { |
||||
|
if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) { |
||||
|
throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType()); |
||||
|
} |
||||
|
List<ListenableFuture<Map.Entry<EntityId, AttributeKvEntry>>> kvFutures = geofencingEntities.stream() |
||||
|
.map(entityId -> { |
||||
|
var attributesFuture = attributesService.find( |
||||
|
tenantId, |
||||
|
entityId, |
||||
|
argument.getRefEntityKey().getScope(), |
||||
|
argument.getRefEntityKey().getKey() |
||||
|
); |
||||
|
return Futures.transform(attributesFuture, resultOpt -> |
||||
|
Map.entry(entityId, resultOpt.orElseGet(() -> |
||||
|
new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), |
||||
|
calculatedFieldCallbackExecutor |
||||
|
); |
||||
|
}).collect(Collectors.toList()); |
||||
|
|
||||
|
ListenableFuture<List<Map.Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
||||
|
|
||||
|
return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream() |
||||
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))), MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected ListenableFuture<ArgumentEntry> fetchArgumentValue(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { |
||||
|
return switch (argument.getRefEntityKey().getType()) { |
||||
|
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument, startTs); |
||||
|
case ATTRIBUTE -> fetchAttribute(tenantId, entityId, argument, startTs); |
||||
|
case TS_LATEST -> fetchTsLatest(tenantId, entityId, argument, startTs); |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { |
||||
|
long argTimeWindow = argument.getTimeWindow() == 0 ? queryEndTs : argument.getTimeWindow(); |
||||
|
long startInterval = queryEndTs - argTimeWindow; |
||||
|
ReadTsKvQuery query = buildTsRollingQuery(tenantId, argument, startInterval, queryEndTs); |
||||
|
|
||||
|
log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query); |
||||
|
ListenableFuture<List<TsKvEntry>> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
||||
|
return Futures.transform(tsRollingFuture, tsRolling -> { |
||||
|
log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, tsRolling == null ? 0 : tsRolling.size(), query); |
||||
|
return ArgumentEntry.createTsRollingArgument(tsRolling, query.getLimit(), argTimeWindow); |
||||
|
}, calculatedFieldCallbackExecutor); |
||||
|
} |
||||
|
|
||||
|
private ListenableFuture<ArgumentEntry> fetchAttribute(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) { |
||||
|
log.trace("[{}][{}] Fetching attribute for key {}", tenantId, entityId, argument.getRefEntityKey()); |
||||
|
var attributeOptFuture = attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()); |
||||
|
|
||||
|
return Futures.transform(attributeOptFuture, attrOpt -> { |
||||
|
log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt); |
||||
|
AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, 0L)); |
||||
|
return transformSingleValueArgument(Optional.of(attributeKvEntry)); |
||||
|
}, calculatedFieldCallbackExecutor); |
||||
|
} |
||||
|
|
||||
|
protected ListenableFuture<ArgumentEntry> fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { |
||||
|
String timeseriesKey = argument.getRefEntityKey().getKey(); |
||||
|
log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, timeseriesKey); |
||||
|
return transformSingleValueArgument( |
||||
|
Futures.transform( |
||||
|
timeseriesService.findLatest(tenantId, entityId, timeseriesKey), |
||||
|
result -> { |
||||
|
log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, timeseriesKey, result); |
||||
|
return result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))); |
||||
|
}, calculatedFieldCallbackExecutor)); |
||||
|
} |
||||
|
|
||||
|
private ReadTsKvQuery buildTsRollingQuery(TenantId tenantId, Argument argument, long startTs, long endTs) { |
||||
|
long maxDataPoints = apiLimitService.getLimit( |
||||
|
tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); |
||||
|
int argumentLimit = argument.getLimit(); |
||||
|
int limit = argumentLimit == 0 || argumentLimit > maxDataPoints ? (int) maxDataPoints : argumentLimit; |
||||
|
return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); |
||||
|
} |
||||
|
|
||||
|
} |
||||