Files
Maries 4589841b0c feat: introduce trigger (#482)
* feat(plugins): add FetchPluginReadme endpoint and update launch configurations

* feat: add PluginReadme database model

* feat: implement readme extracting and storage

* feat: implement readme endpoint

* feat: add plugin asset extraction endpoint with caching support

* Implement trigger functionality and clean up dynamic select code

- Added new trigger-related access types and actions in access.go.
- Introduced new HTTP routes for trigger operations in http_server.gen.go.
- Updated plugin declaration to include triggers in plugin_entities.
- Removed unused dynamic select service and controller files.
- Enhanced message handling in debugging_runtime to support trigger registration.

This update enhances the plugin system by integrating trigger capabilities while cleaning up legacy code.

* Refactor trigger-related types and enhance dynamic select functionality

- Updated TriggerProviderIdentity and TriggerProviderConfiguration to improve structure and validation.
- Renamed TriggerConfiguration to TriggerDeclaration for clarity.
- Added CredentialType to RequestDynamicParameterSelect for better request handling.
- Enhanced PluginDecoderHelper to read and unmarshal trigger files correctly.

These changes streamline the trigger system and improve the overall request handling in the plugin architecture.

* Add trigger functionality and enhance database integration

- Introduced TriggerInstallation model for managing trigger installations in the database.
- Updated autoMigrate function to include trigger installations in the migration process.
- Added new HTTP routes for listing and retrieving triggers in the HTTP server.
- Implemented ListTriggers and GetTrigger controller functions for handling trigger requests.
- Enhanced plugin management functions to create, update, and delete trigger installations during plugin lifecycle events.

These changes integrate trigger capabilities into the system, improving the overall plugin functionality and management.

* feat: add remapping for trigger icons in MediaBucket

- Enhanced the RemapAssets function to include remapping of trigger identity icons and dark icons.
- Added error handling for remapping failures to ensure robust asset management.

These changes improve the handling of trigger assets within the plugin system, ensuring icons are correctly remapped as needed.

* feat: add Multiple field to TriggerParameter for enhanced configuration

- Introduced a new Multiple field in the TriggerParameter struct to allow for multiple values in trigger configurations.
- This addition improves the flexibility of trigger parameters within the plugin system.

These changes enhance the capability of trigger parameters, enabling more complex configurations.

* feat: add Multiple field to ProviderConfig for enhanced configuration

- Introduced a new Multiple field in the ProviderConfig struct to allow for multiple values in provider configurations.
- This addition improves the flexibility of provider options within the plugin system.

These changes enhance the capability of provider configurations, enabling more complex setups.

* fix(plugin): update validation error messages in ManifestValidate method

- Enhanced error messages in the ManifestValidate function to include 'trigger' in the validation checks for plugin declarations.
- Updated logic to ensure that all relevant fields are considered when validating the presence of mutually exclusive parameters.

* feat(trigger): add CHECKBOX parameter type to plugin entities  and refactor the trigger provider strcuture

- Introduced a new CHECKBOX parameter type in constant.go for plugin entities.
- Updated tool_declaration.go and trigger_declaration.go to include TOOL_PARAMETER_TYPE_CHECKBOX and TRIGGER_PARAMETER_TYPE_CHECKBOX respectively.
- Enhanced validation logic to accommodate the new CHECKBOX type in parameter checks.

* fix(trigger): update SubscriptionSchema validation in TriggerProviderDeclaration

- Changed SubscriptionSchema validation from 'omitempty' to 'required' in TriggerProviderDeclaration to ensure it is always provided.
- Updated SubscriptionConstructor field to be a pointer to allow for optional inclusion in the trigger provider configuration.

* fix(trigger): rename ParametersSchema to Parameters in SubscriptionConstructor

- Updated the SubscriptionConstructor struct to rename the ParametersSchema field to Parameters for consistency.
- Adjusted related JSON and YAML marshaling logic to reflect the new field name, ensuring proper handling of trigger parameters.

* refactor(trigger): enhance YAML unmarshalling for SubscriptionConstructor and SubscriptionSchema

- Introduced a new helper function to convert YAML nodes to ProviderConfig lists, improving the handling of subscription_schema and credentials_schema.
- Updated the UnmarshalYAML method to utilize the new function, simplifying the logic for parsing different YAML formats.
- Ensured proper initialization of SubscriptionConstructor fields to prevent nil pointer dereferences.

* fix(trigger): update SubscriptionConstructor validation in TriggerProviderDeclaration

- Changed the validation for SubscriptionConstructor in TriggerProviderDeclaration from 'omitempty,dive' to 'omitempty' to simplify the validation logic.
- Ensured that the SubscriptionConstructor field remains optional while maintaining its intended functionality.

* refactor(trigger): rename Trigger to Event in plugin entities and related structures

- Updated the naming conventions in trigger_declaration.go to replace 'Trigger' with 'Event' for better clarity and consistency.
- Adjusted related types, validation functions, and unmarshalling logic to reflect the new 'Event' terminology.
- Ensured that all references to triggers in the codebase are updated to events, including in the SubscriptionConstructor and response structures.

* refactor(trigger): rename TriggerInvoke to TriggerInvokeEvent and update related structures

- Renamed TriggerInvoke function and associated request/response types to TriggerInvokeEvent for improved clarity.
- Updated routing and controller methods to reflect the new naming convention.
- Ensured all references to the trigger invoke functionality are consistent with the new event terminology.

* refactor(trigger): remove Subscription struct from trigger_declaration.go and update TriggerDispatchEventRequest

- Removed the Subscription struct from trigger_declaration.go to streamline the codebase.
- Added Credentials field to TriggerDispatchEventRequest for enhanced functionality and clarity.
- Ensured that the changes maintain consistency with the existing naming conventions and structures.

* fix(trigger): improve nil checks for SubscriptionConstructor in TriggerProviderDeclaration

- Added nil checks for SubscriptionConstructor before accessing its fields to prevent potential nil pointer dereferences.
- Ensured that Parameters and CredentialsSchema are initialized only if SubscriptionConstructor is not nil, enhancing code robustness.

* fix(plugin): add recovery mechanism in OnTraffic to handle panics

- Introduced a deferred function in OnTraffic to recover from panics, logging the error and stack trace for better debugging.
- This enhancement improves the stability of the DifyServer by preventing crashes due to unexpected runtime errors.

* feat(trigger): add Subscription field to TriggerInvokeEventRequest

- Introduced a new Subscription field in the TriggerInvokeEventRequest struct to accommodate subscription data.
- Ensured the field is marked as required, enhancing the request's functionality and validation requirements.

* refactor(event): simplify EventDescription structure in EventDeclaration

- Removed the EventDescription struct and replaced it with a direct I18nObject field in EventDeclaration.
- This change streamlines the event configuration by reducing complexity while maintaining required validation for the description.

* feat(trigger): add UserID field to TriggerDispatchEventResponse

- Introduced a new UserID field in the TriggerDispatchEventResponse struct to include user identification in the response.
- The field is marked as optional, enhancing the response's flexibility while maintaining existing functionality.

* feat: add payload to TriggerDispatchEventResponse

* fix

* feat(trigger): update TriggerDispatchEventResponse structure

* fix: avoid path collusion

* fix: missing )

* fix: query param

* fix: form param

* fix: remove redundant dynamic parameter access type

* fix: remove dynamic parameter access type from validation

* Update internal/server/controllers/plugins.go

Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>

---------

Co-authored-by: Stream <Stream_2@qq.com>
Co-authored-by: Yeuoly <admin@srmxy.cn>
Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com>
2025-10-29 13:56:58 +08:00

304 lines
11 KiB
Go

package controllers
import (
"errors"
"net/http"
"strings"
"github.com/gin-gonic/gin"
"github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager"
"github.com/langgenius/dify-plugin-daemon/internal/service"
"github.com/langgenius/dify-plugin-daemon/internal/types/app"
"github.com/langgenius/dify-plugin-daemon/internal/types/exception"
"github.com/langgenius/dify-plugin-daemon/pkg/entities/constants"
"github.com/langgenius/dify-plugin-daemon/pkg/entities/plugin_entities"
)
func GetAsset(c *gin.Context) {
pluginManager := plugin_manager.Manager()
asset, err := pluginManager.GetAsset(c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, exception.InternalServerError(err).ToResponse())
return
}
c.Data(http.StatusOK, "application/octet-stream", asset)
}
func UploadPlugin(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
difyPkgFileHeader, err := c.FormFile("dify_pkg")
if err != nil {
c.JSON(http.StatusOK, exception.BadRequestError(err).ToResponse())
return
}
tenantId := c.Param("tenant_id")
if tenantId == "" {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("tenant ID is required")).ToResponse())
return
}
if difyPkgFileHeader.Size > app.MaxPluginPackageSize {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("file size exceeds the maximum limit")).ToResponse())
return
}
verifySignature := c.PostForm("verify_signature") == "true"
difyPkgFile, err := difyPkgFileHeader.Open()
if err != nil {
c.JSON(http.StatusOK, exception.BadRequestError(err).ToResponse())
return
}
defer difyPkgFile.Close()
c.JSON(http.StatusOK, service.UploadPluginPkg(app, c, tenantId, difyPkgFile, verifySignature))
}
}
func UploadBundle(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
difyBundleFileHeader, err := c.FormFile("dify_bundle")
if err != nil {
c.JSON(http.StatusOK, exception.BadRequestError(err).ToResponse())
return
}
tenantId := c.Param("tenant_id")
if tenantId == "" {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("tenant ID is required")).ToResponse())
return
}
if difyBundleFileHeader.Size > app.MaxBundlePackageSize {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("file size exceeds the maximum limit")).ToResponse())
return
}
verifySignature := c.PostForm("verify_signature") == "true"
difyBundleFile, err := difyBundleFileHeader.Open()
if err != nil {
c.JSON(http.StatusOK, exception.BadRequestError(err).ToResponse())
return
}
defer difyBundleFile.Close()
c.JSON(http.StatusOK, service.UploadPluginBundle(app, c, tenantId, difyBundleFile, verifySignature))
}
}
func UpgradePlugin(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
OriginalPluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `json:"original_plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
NewPluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `json:"new_plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
Source string `json:"source" validate:"required"`
Meta map[string]any `json:"meta" validate:"omitempty"`
}) {
if request.TenantID == constants.GlobalTenantId && !app.PluginAllowOrphans {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("orphan plugin is not allowed")).ToResponse())
return
}
c.JSON(http.StatusOK, service.UpgradePlugin(
app,
request.TenantID,
request.Source,
request.Meta,
request.OriginalPluginUniqueIdentifier,
request.NewPluginUniqueIdentifier,
))
})
}
}
func InstallPluginFromIdentifiers(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginUniqueIdentifiers []plugin_entities.PluginUniqueIdentifier `json:"plugin_unique_identifiers" validate:"required,max=64,dive,plugin_unique_identifier"`
Source string `json:"source" validate:"required"`
Metas []map[string]any `json:"metas" validate:"omitempty"`
}) {
if request.Metas == nil {
request.Metas = []map[string]any{}
}
if request.TenantID == constants.GlobalTenantId && !app.PluginAllowOrphans {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("orphan plugin is not allowed")).ToResponse())
return
}
if len(request.Metas) != len(request.PluginUniqueIdentifiers) {
c.JSON(http.StatusOK, exception.BadRequestError(errors.New("the number of metas must be equal to the number of plugin unique identifiers")).ToResponse())
return
}
for i := range request.Metas {
if request.Metas[i] == nil {
request.Metas[i] = map[string]any{}
}
}
c.JSON(http.StatusOK, service.InstallPluginFromIdentifiers(
app, request.TenantID, request.PluginUniqueIdentifiers, request.Source, request.Metas,
))
})
}
}
func ReinstallPluginFromIdentifier(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
BindRequest(c, func(request struct {
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `json:"plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
}) {
service.ReinstallPluginFromIdentifier(c, app, request.PluginUniqueIdentifier)
})
}
}
func DecodePluginFromIdentifier(app *app.Config) gin.HandlerFunc {
return func(c *gin.Context) {
BindRequest(c, func(request struct {
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `json:"plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
}) {
c.JSON(http.StatusOK, service.DecodePluginFromIdentifier(app, request.PluginUniqueIdentifier))
})
}
}
func FetchPluginInstallationTasks(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
Page int `form:"page" validate:"required,min=1"`
PageSize int `form:"page_size" validate:"required,min=1,max=256"`
}) {
c.JSON(http.StatusOK, service.FetchPluginInstallationTasks(request.TenantID, request.Page, request.PageSize))
})
}
func FetchPluginInstallationTask(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
TaskID string `uri:"id" validate:"required"`
}) {
c.JSON(http.StatusOK, service.FetchPluginInstallationTask(request.TenantID, request.TaskID))
})
}
func DeletePluginInstallationTask(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
TaskID string `uri:"id" validate:"required"`
}) {
c.JSON(http.StatusOK, service.DeletePluginInstallationTask(request.TenantID, request.TaskID))
})
}
func DeleteAllPluginInstallationTasks(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
}) {
c.JSON(http.StatusOK, service.DeleteAllPluginInstallationTasks(request.TenantID))
})
}
func DeletePluginInstallationItemFromTask(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
TaskID string `uri:"id" validate:"required"`
Identifier string `uri:"identifier" validate:"required"`
}) {
identifierString := strings.TrimLeft(request.Identifier, "/")
identifier, err := plugin_entities.NewPluginUniqueIdentifier(identifierString)
if err != nil {
c.JSON(http.StatusOK, exception.BadRequestError(err).ToResponse())
return
}
c.JSON(http.StatusOK, service.DeletePluginInstallationItemFromTask(request.TenantID, request.TaskID, identifier))
})
}
func FetchPluginManifest(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `form:"plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
}) {
c.JSON(http.StatusOK, service.FetchPluginManifest(request.PluginUniqueIdentifier))
})
}
func FetchPluginReadme(c *gin.Context) {
BindRequest(c, func(request struct {
TenantId string `uri:"tenant_id" validate:"required"`
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `form:"plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
Language string `form:"language" validate:"omitempty"`
}) {
c.JSON(http.StatusOK, service.FetchPluginReadme(request.TenantId, request.PluginUniqueIdentifier, request.Language))
})
}
func UninstallPlugin(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginInstallationID string `json:"plugin_installation_id" validate:"required"`
}) {
c.JSON(http.StatusOK, service.UninstallPlugin(request.TenantID, request.PluginInstallationID))
})
}
func FetchPluginFromIdentifier(c *gin.Context) {
BindRequest(c, func(request struct {
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier `form:"plugin_unique_identifier" validate:"required,plugin_unique_identifier"`
}) {
c.JSON(http.StatusOK, service.FetchPluginFromIdentifier(request.PluginUniqueIdentifier))
})
}
func ListPlugins(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
Page int `form:"page" validate:"required,min=1"`
PageSize int `form:"page_size" validate:"required,min=1,max=256"`
}) {
c.JSON(http.StatusOK, service.ListPlugins(request.TenantID, request.Page, request.PageSize))
})
}
func BatchFetchPluginInstallationByIDs(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginIDs []string `json:"plugin_ids" validate:"required,max=256"`
}) {
c.JSON(http.StatusOK, service.BatchFetchPluginInstallationByIDs(request.TenantID, request.PluginIDs))
})
}
func FetchMissingPluginInstallations(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginUniqueIdentifiers []plugin_entities.PluginUniqueIdentifier `json:"plugin_unique_identifiers" validate:"required,max=256,dive,plugin_unique_identifier"`
}) {
c.JSON(http.StatusOK, service.FetchMissingPluginInstallations(request.TenantID, request.PluginUniqueIdentifiers))
})
}
func ExtractPluginAsset(c *gin.Context) {
BindRequest(c, func(request struct {
TenantID string `uri:"tenant_id" validate:"required"`
PluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier "form:\"plugin_unique_identifier\" validate:\"required,plugin_unique_identifier\""
FilePath string `form:"file_path" validate:"required"`
}) {
manager := plugin_manager.Manager()
asset, err := manager.ExtractPluginAsset(request.PluginUniqueIdentifier, request.FilePath)
if err != nil {
c.JSON(http.StatusInternalServerError, exception.InternalServerError(err).ToResponse())
return
}
c.Data(http.StatusOK, "application/octet-stream", asset)
})
}