Files
dify-plugin-daemon/internal/service/manage_plugin.go
Yeuoly 878edde455 feat(datasource): Implement datasource validation and invocation steps (#295)
* feat(datasource): Implement datasource validation and invocation steps

- Added new functionality for validating datasource credentials and invoking the first and second steps of the datasource process.
- Introduced new API endpoints for datasource operations in the HTTP server.
- Created corresponding service and controller methods to handle datasource requests.
- Defined new request and response structures for datasource operations in the entities package.

* feat: add routes

* feat(datasource): Add initial datasource models and validation

- Introduced new `DatasourceInstallation` model to represent datasource installations with relevant fields.
- Created `datasource.go` file for future datasource service implementation.
- Enhanced `datasource_declaration.go` with new types and validation functions for datasource provider and parameters.
- Updated `plugin_declaration.go` to include datasource support in plugin structures.
- Refactored `agent_declaration.go` and `tool_declaration.go` to use a unified `ParameterOption` type for options.

* refactor(datasource): Update request and response types for online document content retrieval

- Renamed and updated request and response types for the `DatasourceGetOnlineDocumentPageContent` function to improve clarity and consistency.
- Introduced `RequestInvokeOnlineDocumentDatasourceGetContent` and `DatasourceInvokeOnlineDocumentGetContentResponse` types.
- Adjusted related function signatures and dispatchers to reflect the new types across the datasource service and controller implementations.

* feat(datasource): Implement datasource installation handling

- Added functionality to create and update `DatasourceInstallation` records during plugin installation and upgrade processes.
- Enhanced the `InstallPlugin` function to create a new datasource installation if a datasource declaration is present.
- Updated the `UpgradePlugin` function to handle the deletion of the original datasource installation and creation of a new one if the datasource declaration changes.

* feat(datasource): Add datasource registration handling in plugin runtime

- Introduced handling for datasource declarations in the plugin runtime.
- Updated `RemotePluginRuntime` to track datasource registration status.
- Enhanced message processing to include validation and assignment of datasource declarations.

* feat(datasource): Implement datasource listing and retrieval endpoints

- Added `ListDatasources` and `GetDatasource` functions to the service layer for handling datasource queries.
- Created corresponding controller methods to process HTTP requests for listing and retrieving datasources.
- Implemented request validation for tenant ID, page, page size, plugin ID, and provider parameters.

* add datasource routes

* feat(datasource): Add icon remapping for datasource declarations

- Implemented functionality to remap icons for both the main datasource and its individual datasources within the plugin declaration.
- Enhanced error handling to provide clearer feedback when remapping fails.

* fix(datasource): update OAuthSchema validation to remove unnecessary 'dive' tag

- Modified the validation tag for OAuthSchema in DatasourceProviderDeclaration to simplify the validation process by removing the 'dive' requirement.

* feat(tests): add datasource declaration parsing and output

- Introduced a new main.go file for testing datasource declaration parsing.
- Implemented JSON unmarshalling for RemotePluginRegisterPayload and DatasourceProviderDeclaration.
- Added a sample datasource declaration for testing purposes.

* fix(plugin): improve error handling in UninstallPlugin and add datasource deletion

- Enhanced error handling to return a specific message when a plugin is not installed.
- Added functionality to delete the associated datasource installation during the uninstallation process.

* feat(datasource): add support for decoding datasource provider declaration

- Added custom JSON marshalling and unmarshalling methods to handle CredentialsSchema and Datasources more effectively.
- Improved error handling during YAML unmarshalling to support both object and array formats for CredentialsSchema.
- Ensured proper initialization of DatasourceFiles and Tags to prevent nil references.

* fix: provider type

* refactor(datasource): simplify JSON unmarshalling for CredentialsSchema

- Removed complex handling of CredentialsSchema in the UnmarshalJSON method, focusing on the Datasources field.
- Streamlined the code to improve readability and maintainability by eliminating unnecessary checks and logic related to CredentialsSchema.

* feat: streaming datasource

* feat: datasource

* feat: datasource

* feat:datasource

* feat:datasource

* feat:datasource

* feat: add redirect_uri field to OAuth request structs

* feat: add online driver file request and response structures

* feat: add online driver file request and response structures

* feat: add online_driver datasource type to validation

* feat: rename online driver to online drive and update related classes and methods :)

* refactor: rename OnlineDocumentPageChunk to DatasourceGetPagesResponse and update related references

* feat: update request types for online drive browsing and downloading

* feat: add metadata field to OAuthGetCredentialsResult

* feat: built-in json schema definations

* feat(plugin_entities): add built-in schema definitions and processing for datasource YAML

* test(plugin_entities): add unit tests for schema definitions and YAML processing

* refactor(plugin_entities): centralize built-in schema definitions and processing

* refactor(plugin_entities): remove unused properties from built-in schema definitions

* refactor(plugin_entities): built-in schema & new datasource structure

* refactor(plugin_entities): enhance schema processing with  checks and error handling

* refactor(plugin_entities): update validation for OnlineDriveBrowseFilesRequest prefix field to be optional

* refactor(json_schema): remove json schema definitions and validation

* fix(plugin_entities): remove output_schema validation tests

* feat: add rag tag

---------

Co-authored-by: Harry <xh001x@hotmail.com>
Co-authored-by: Dongyu Li <544104925@qq.com>
Co-authored-by: Novice <novice12185727@gmail.com>
2025-09-16 17:32:14 +08:00

576 lines
17 KiB
Go

package service
import (
"errors"
"time"
"github.com/langgenius/dify-plugin-daemon/internal/db"
"github.com/langgenius/dify-plugin-daemon/internal/types/exception"
"github.com/langgenius/dify-plugin-daemon/internal/types/models"
"github.com/langgenius/dify-plugin-daemon/internal/utils/cache/helper"
"github.com/langgenius/dify-plugin-daemon/internal/utils/strings"
"github.com/langgenius/dify-plugin-daemon/pkg/entities"
"github.com/langgenius/dify-plugin-daemon/pkg/entities/manifest_entities"
"github.com/langgenius/dify-plugin-daemon/pkg/entities/plugin_entities"
)
func ListPlugins(tenant_id string, page int, page_size int) *entities.Response {
type installation struct {
ID string `json:"id"`
Name string `json:"name"`
PluginID string `json:"plugin_id"`
TenantID string `json:"tenant_id"`
PluginUniqueIdentifier string `json:"plugin_unique_identifier"`
EndpointsActive int `json:"endpoints_active"`
EndpointsSetups int `json:"endpoints_setups"`
InstallationID string `json:"installation_id"`
Declaration *plugin_entities.PluginDeclaration `json:"declaration"`
RuntimeType plugin_entities.PluginRuntimeType `json:"runtime_type"`
Version manifest_entities.Version `json:"version"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
Source string `json:"source"`
Checksum string `json:"checksum"`
Meta map[string]any `json:"meta"`
}
type responseData struct {
List []installation `json:"list"`
Total int64 `json:"total"`
}
// get total count
totalCount, err := db.GetCount[models.PluginInstallation](
db.Equal("tenant_id", tenant_id),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
pluginInstallations, err := db.GetAll[models.PluginInstallation](
db.Equal("tenant_id", tenant_id),
db.OrderBy("created_at", true),
db.Page(page, page_size),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]installation, 0, len(pluginInstallations))
for _, plugin_installation := range pluginInstallations {
pluginUniqueIdentifier, err := plugin_entities.NewPluginUniqueIdentifier(
plugin_installation.PluginUniqueIdentifier,
)
if err != nil {
return exception.UniqueIdentifierError(err).ToResponse()
}
pluginDeclaration, err := helper.CombinedGetPluginDeclaration(
pluginUniqueIdentifier,
plugin_entities.PluginRuntimeType(plugin_installation.RuntimeType),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data = append(data, installation{
ID: plugin_installation.ID,
Name: pluginDeclaration.Name,
TenantID: plugin_installation.TenantID,
PluginID: pluginUniqueIdentifier.PluginID(),
PluginUniqueIdentifier: pluginUniqueIdentifier.String(),
InstallationID: plugin_installation.ID,
Declaration: pluginDeclaration,
EndpointsSetups: plugin_installation.EndpointsSetups,
EndpointsActive: plugin_installation.EndpointsActive,
RuntimeType: plugin_entities.PluginRuntimeType(plugin_installation.RuntimeType),
Version: pluginDeclaration.Version,
CreatedAt: plugin_installation.CreatedAt,
UpdatedAt: plugin_installation.UpdatedAt,
Source: plugin_installation.Source,
Meta: plugin_installation.Meta,
Checksum: pluginUniqueIdentifier.Checksum(),
})
}
finalData := responseData{
List: data,
Total: totalCount,
}
return entities.NewSuccessResponse(finalData)
}
// Using plugin_ids to fetch plugin installations
func BatchFetchPluginInstallationByIDs(tenant_id string, plugin_ids []string) *entities.Response {
type installation struct {
models.PluginInstallation
Version manifest_entities.Version `json:"version"`
Checksum string `json:"checksum"`
Declaration *plugin_entities.PluginDeclaration `json:"declaration"`
}
if len(plugin_ids) == 0 {
return entities.NewSuccessResponse([]installation{})
}
pluginInstallations, err := db.GetAll[models.PluginInstallation](
db.Equal("tenant_id", tenant_id),
db.InArray("plugin_id", strings.Map(plugin_ids, func(id string) any { return id })),
db.Page(1, 256), // TODO: pagination
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]installation, 0, len(pluginInstallations))
for _, plugin_installation := range pluginInstallations {
pluginUniqueIdentifier, err := plugin_entities.NewPluginUniqueIdentifier(
plugin_installation.PluginUniqueIdentifier,
)
if err != nil {
return exception.InternalServerError(errors.Join(errors.New("invalid plugin unique identifier found"), err)).ToResponse()
}
pluginDeclaration, err := helper.CombinedGetPluginDeclaration(
pluginUniqueIdentifier,
plugin_entities.PluginRuntimeType(plugin_installation.RuntimeType),
)
if err != nil {
return exception.InternalServerError(errors.Join(errors.New("failed to get plugin declaration"), err)).ToResponse()
}
data = append(data, installation{
PluginInstallation: plugin_installation,
Version: pluginUniqueIdentifier.Version(),
Checksum: pluginUniqueIdentifier.Checksum(),
Declaration: pluginDeclaration,
})
}
return entities.NewSuccessResponse(data)
}
// check which plugin is missing
func FetchMissingPluginInstallations(tenant_id string, plugin_unique_identifiers []plugin_entities.PluginUniqueIdentifier) *entities.Response {
type MissingPluginDependency struct {
PluginUniqueIdentifier string `json:"plugin_unique_identifier"`
CurrentIdentifier string `json:"current_identifier"`
}
result := make([]MissingPluginDependency, 0, len(plugin_unique_identifiers))
if len(plugin_unique_identifiers) == 0 {
return entities.NewSuccessResponse(result)
}
installed, err := db.GetAll[models.PluginInstallation](
db.Equal("tenant_id", tenant_id),
db.InArray(
"plugin_id",
strings.Map(
plugin_unique_identifiers,
func(id plugin_entities.PluginUniqueIdentifier) any {
return id.PluginID()
},
),
),
db.Page(1, 256), // TODO: pagination
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
// check which plugin is missing
for _, pluginUniqueIdentifier := range plugin_unique_identifiers {
found := false
for _, installedPlugin := range installed {
if installedPlugin.PluginID == pluginUniqueIdentifier.PluginID() {
found = true
if installedPlugin.PluginUniqueIdentifier != pluginUniqueIdentifier.String() {
// version mismatched
result = append(result, MissingPluginDependency{
PluginUniqueIdentifier: pluginUniqueIdentifier.String(),
CurrentIdentifier: installedPlugin.PluginUniqueIdentifier,
})
}
break
}
}
if !found {
result = append(result, MissingPluginDependency{
PluginUniqueIdentifier: pluginUniqueIdentifier.String(),
})
}
}
return entities.NewSuccessResponse(result)
}
func ListTools(tenant_id string, page int, page_size int) *entities.Response {
type Tool struct {
models.ToolInstallation // pointer to avoid deep copy
Declaration *plugin_entities.ToolProviderDeclaration `json:"declaration"`
}
providers, err := db.GetAll[models.ToolInstallation](
db.Equal("tenant_id", tenant_id),
db.Page(page, page_size),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]Tool, 0, len(providers))
for _, provider := range providers {
// check if plugin id starts with uuid
// split by uuid length
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(provider.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data = append(data, Tool{
ToolInstallation: provider,
Declaration: declaration.Tool,
})
}
return entities.NewSuccessResponse(data)
}
func ListModels(tenant_id string, page int, page_size int) *entities.Response {
type AIModel struct {
models.AIModelInstallation // pointer to avoid deep copy
Declaration *plugin_entities.ModelProviderDeclaration `json:"declaration"`
}
providers, err := db.GetAll[models.AIModelInstallation](
db.Equal("tenant_id", tenant_id),
db.Page(page, page_size),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]AIModel, 0, len(providers))
for _, provider := range providers {
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(provider.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data = append(data, AIModel{
AIModelInstallation: provider,
Declaration: declaration.Model,
})
}
return entities.NewSuccessResponse(data)
}
func GetTool(tenant_id string, plugin_id string, provider string) *entities.Response {
type Tool struct {
models.ToolInstallation // pointer to avoid deep copy
Declaration *plugin_entities.ToolProviderDeclaration `json:"declaration"`
}
// try get tool
tool, err := db.GetOne[models.ToolInstallation](
db.Equal("tenant_id", tenant_id),
db.Equal("plugin_id", plugin_id),
)
if err != nil {
if err == db.ErrDatabaseNotFound {
return exception.ErrPluginNotFound().ToResponse()
}
return exception.InternalServerError(err).ToResponse()
}
if tool.Provider != provider {
return exception.ErrPluginNotFound().ToResponse()
}
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(tool.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
return entities.NewSuccessResponse(Tool{
ToolInstallation: tool,
Declaration: declaration.Tool,
})
}
type RequestCheckToolExistence struct {
PluginID string `json:"plugin_id" validate:"required"`
ProviderName string `json:"provider_name" validate:"required"`
}
func CheckToolExistence(tenantId string, providerIds []RequestCheckToolExistence) *entities.Response {
existence := make([]bool, 0, len(providerIds))
// get all providers
providers, err := db.GetAll[models.ToolInstallation](
db.Equal("tenant_id", tenantId),
db.InArray("plugin_id", strings.Map(providerIds, func(id RequestCheckToolExistence) any { return id.PluginID })),
db.Page(1, 256), // TODO: pagination
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
// check provider id
for _, providerId := range providerIds {
found := false
for _, provider := range providers {
if provider.PluginID == providerId.PluginID && provider.Provider == providerId.ProviderName {
found = true
break
}
}
existence = append(existence, found)
}
return entities.NewSuccessResponse(existence)
}
func ListAgentStrategies(tenant_id string, page int, page_size int) *entities.Response {
type AgentStrategy struct {
models.AgentStrategyInstallation // pointer to avoid deep copy
Declaration *plugin_entities.AgentStrategyProviderDeclaration `json:"declaration"`
Meta plugin_entities.PluginMeta `json:"meta"`
}
providers, err := db.GetAll[models.AgentStrategyInstallation](
db.Equal("tenant_id", tenant_id),
db.Page(page, page_size),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]AgentStrategy, 0, len(providers))
for _, provider := range providers {
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(provider.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data = append(data, AgentStrategy{
AgentStrategyInstallation: provider,
Declaration: declaration.AgentStrategy,
Meta: declaration.Meta,
})
}
return entities.NewSuccessResponse(data)
}
func GetAgentStrategy(tenant_id string, plugin_id string, provider string) *entities.Response {
type AgentStrategy struct {
models.AgentStrategyInstallation // pointer to avoid deep copy
Declaration *plugin_entities.AgentStrategyProviderDeclaration `json:"declaration"`
Meta plugin_entities.PluginMeta `json:"meta"`
}
agent_strategy, err := db.GetOne[models.AgentStrategyInstallation](
db.Equal("tenant_id", tenant_id),
db.Equal("plugin_id", plugin_id),
)
if err != nil {
if err == db.ErrDatabaseNotFound {
return exception.ErrPluginNotFound().ToResponse()
}
return exception.InternalServerError(err).ToResponse()
}
if agent_strategy.Provider != provider {
return exception.ErrPluginNotFound().ToResponse()
}
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(agent_strategy.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
return entities.NewSuccessResponse(AgentStrategy{
AgentStrategyInstallation: agent_strategy,
Declaration: declaration.AgentStrategy,
Meta: declaration.Meta,
})
}
func ListDatasources(tenant_id string, page int, page_size int) *entities.Response {
type Datasource struct {
models.DatasourceInstallation // pointer to avoid deep copy
Declaration *plugin_entities.DatasourceProviderDeclaration `json:"declaration"`
}
providers, err := db.GetAll[models.DatasourceInstallation](
db.Equal("tenant_id", tenant_id),
db.Page(page, page_size),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data := make([]Datasource, 0, len(providers))
for _, provider := range providers {
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(provider.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
data = append(data, Datasource{
DatasourceInstallation: provider,
Declaration: declaration.Datasource,
})
}
return entities.NewSuccessResponse(data)
}
func GetDatasource(tenant_id string, plugin_id string, provider string) *entities.Response {
type Datasource struct {
models.DatasourceInstallation // pointer to avoid deep copy
Declaration *plugin_entities.DatasourceProviderDeclaration `json:"declaration"`
}
datasource, err := db.GetOne[models.DatasourceInstallation](
db.Equal("tenant_id", tenant_id),
db.Equal("plugin_id", plugin_id),
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
if datasource.Provider != provider {
return exception.ErrPluginNotFound().ToResponse()
}
uniqueIdentifier := plugin_entities.PluginUniqueIdentifier(datasource.PluginUniqueIdentifier)
var runtimeType plugin_entities.PluginRuntimeType
if uniqueIdentifier.RemoteLike() {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_REMOTE
} else {
runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
}
declaration, err := helper.CombinedGetPluginDeclaration(
uniqueIdentifier,
runtimeType,
)
if err != nil {
return exception.InternalServerError(err).ToResponse()
}
return entities.NewSuccessResponse(Datasource{
DatasourceInstallation: datasource,
Declaration: declaration.Datasource,
})
}