mirror of
https://github.com/langgenius/dify-plugin-daemon.git
synced 2026-07-23 10:15:22 -04:00
4589841b0c
* 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>
307 lines
8.9 KiB
Go
307 lines
8.9 KiB
Go
package plugin_manager
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
|
|
lru "github.com/hashicorp/golang-lru/v2"
|
|
"github.com/langgenius/dify-cloud-kit/oss"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/dify_invocation"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/dify_invocation/real"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager/debugging_runtime"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager/media_transport"
|
|
serverless "github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager/serverless_connector"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/db"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/types/app"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/types/models"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/cache"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/cache/helper"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/lock"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/log"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/mapping"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/entities/plugin_entities"
|
|
"github.com/langgenius/dify-plugin-daemon/pkg/plugin_packager/decoder"
|
|
)
|
|
|
|
type PluginManager struct {
|
|
m mapping.Map[string, plugin_entities.PluginLifetime]
|
|
|
|
// mediaBucket is used to manage media files like plugin icons, images, etc.
|
|
mediaBucket *media_transport.MediaBucket
|
|
|
|
// packageBucket is used to manage plugin packages, all the packages uploaded by users will be saved here
|
|
packageBucket *media_transport.PackageBucket
|
|
|
|
// installedBucket is used to manage installed plugins, all the installed plugins will be saved here
|
|
installedBucket *media_transport.InstalledBucket
|
|
|
|
// register plugin
|
|
pluginRegisters []func(lifetime plugin_entities.PluginLifetime) error
|
|
|
|
// localPluginLaunchingLock is a lock to launch local plugins
|
|
localPluginLaunchingLock *lock.GranularityLock
|
|
|
|
// backwardsInvocation is a handle to invoke dify
|
|
backwardsInvocation dify_invocation.BackwardsInvocation
|
|
|
|
config *app.Config
|
|
|
|
// remote plugin server
|
|
remotePluginServer debugging_runtime.RemotePluginServerInterface
|
|
|
|
// max launching lock to prevent too many plugins launching at the same time
|
|
maxLaunchingLock chan bool
|
|
}
|
|
|
|
var (
|
|
manager *PluginManager
|
|
)
|
|
|
|
func InitGlobalManager(oss oss.OSS, configuration *app.Config) *PluginManager {
|
|
manager = &PluginManager{
|
|
mediaBucket: media_transport.NewAssetsBucket(
|
|
oss,
|
|
configuration.PluginMediaCachePath,
|
|
configuration.PluginMediaCacheSize,
|
|
),
|
|
packageBucket: media_transport.NewPackageBucket(
|
|
oss,
|
|
configuration.PluginPackageCachePath,
|
|
),
|
|
installedBucket: media_transport.NewInstalledBucket(
|
|
oss,
|
|
configuration.PluginInstalledPath,
|
|
),
|
|
localPluginLaunchingLock: lock.NewGranularityLock(),
|
|
// By default, we allow up to configuration.PluginLocalLaunchingConcurrent plugins to be launched concurrently; if not configured, the default is 2.
|
|
maxLaunchingLock: make(chan bool, configuration.PluginLocalLaunchingConcurrent),
|
|
config: configuration,
|
|
}
|
|
|
|
// initialize plugin asset cache
|
|
if pluginAssetCache == nil {
|
|
c, err := lru.New[string, []byte](256)
|
|
if err != nil {
|
|
log.Panic("init plugin asset cache failed: %s", err.Error())
|
|
}
|
|
pluginAssetCache = c
|
|
}
|
|
|
|
return manager
|
|
}
|
|
|
|
func Manager() *PluginManager {
|
|
return manager
|
|
}
|
|
|
|
func (p *PluginManager) Get(
|
|
identity plugin_entities.PluginUniqueIdentifier,
|
|
) (plugin_entities.PluginLifetime, error) {
|
|
if identity.RemoteLike() || p.config.Platform == app.PLATFORM_LOCAL {
|
|
// check if it's a debugging plugin or a local plugin
|
|
if v, ok := p.m.Load(identity.String()); ok {
|
|
return v, nil
|
|
}
|
|
return nil, errors.New("plugin not found")
|
|
} else {
|
|
// otherwise, use serverless runtime instead
|
|
pluginSessionInterface, err := p.getServerlessPluginRuntime(identity)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return pluginSessionInterface, nil
|
|
}
|
|
}
|
|
|
|
func (p *PluginManager) GetAsset(id string) ([]byte, error) {
|
|
return p.mediaBucket.Get(id)
|
|
}
|
|
|
|
func (p *PluginManager) Launch(configuration *app.Config) {
|
|
log.Info("start plugin manager daemon...")
|
|
|
|
// init redis client
|
|
if configuration.RedisUseSentinel {
|
|
// use Redis Sentinel
|
|
sentinels := strings.Split(configuration.RedisSentinels, ",")
|
|
if err := cache.InitRedisSentinelClient(
|
|
sentinels,
|
|
configuration.RedisSentinelServiceName,
|
|
configuration.RedisUser,
|
|
configuration.RedisPass,
|
|
configuration.RedisSentinelUsername,
|
|
configuration.RedisSentinelPassword,
|
|
configuration.RedisUseSsl,
|
|
configuration.RedisDB,
|
|
configuration.RedisSentinelSocketTimeout,
|
|
); err != nil {
|
|
log.Panic("init redis sentinel client failed: %s", err.Error())
|
|
}
|
|
} else {
|
|
if err := cache.InitRedisClient(
|
|
fmt.Sprintf("%s:%d", configuration.RedisHost, configuration.RedisPort),
|
|
configuration.RedisUser,
|
|
configuration.RedisPass,
|
|
configuration.RedisUseSsl,
|
|
configuration.RedisDB,
|
|
); err != nil {
|
|
log.Panic("init redis client failed: %s", err.Error())
|
|
}
|
|
}
|
|
|
|
invocation, err := real.NewDifyInvocationDaemon(
|
|
real.NewDifyInvocationDaemonPayload{
|
|
BaseUrl: configuration.DifyInnerApiURL,
|
|
CallingKey: configuration.DifyInnerApiKey,
|
|
WriteTimeout: configuration.DifyInvocationWriteTimeout,
|
|
ReadTimeout: configuration.DifyInvocationReadTimeout,
|
|
},
|
|
)
|
|
if err != nil {
|
|
log.Panic("init dify invocation daemon failed: %s", err.Error())
|
|
}
|
|
p.backwardsInvocation = invocation
|
|
|
|
// start local watcher
|
|
if configuration.Platform == app.PLATFORM_LOCAL {
|
|
p.startLocalWatcher(configuration)
|
|
}
|
|
|
|
// launch serverless connector
|
|
if configuration.Platform == app.PLATFORM_SERVERLESS {
|
|
serverless.Init(configuration)
|
|
}
|
|
|
|
// start remote watcher
|
|
p.startRemoteWatcher(configuration)
|
|
}
|
|
|
|
func (p *PluginManager) BackwardsInvocation() dify_invocation.BackwardsInvocation {
|
|
return p.backwardsInvocation
|
|
}
|
|
|
|
func (p *PluginManager) SavePackage(plugin_unique_identifier plugin_entities.PluginUniqueIdentifier, pkg []byte, thirdPartySignatureVerificationConfig *decoder.ThirdPartySignatureVerificationConfig) (
|
|
*plugin_entities.PluginDeclaration, error,
|
|
) {
|
|
// try to decode the package
|
|
packageDecoder, err := decoder.NewZipPluginDecoderWithThirdPartySignatureVerificationConfig(pkg, thirdPartySignatureVerificationConfig)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// get the declaration
|
|
declaration, err := packageDecoder.Manifest()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if err := declaration.ManifestValidate(); err != nil {
|
|
return nil, errors.Join(err, fmt.Errorf("illegal plugin manifest"))
|
|
}
|
|
|
|
// get the assets
|
|
assets, err := packageDecoder.Assets()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// remap the assets
|
|
_, err = p.mediaBucket.RemapAssets(&declaration, assets)
|
|
if err != nil {
|
|
return nil, errors.Join(err, fmt.Errorf("failed to remap assets"))
|
|
}
|
|
|
|
uniqueIdentifier, err := packageDecoder.UniqueIdentity()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// save to storage
|
|
err = p.packageBucket.Save(plugin_unique_identifier.String(), pkg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// create plugin if not exists
|
|
if _, err := db.GetOne[models.PluginDeclaration](
|
|
db.Equal("plugin_unique_identifier", uniqueIdentifier.String()),
|
|
); err == db.ErrDatabaseNotFound {
|
|
err = db.Create(&models.PluginDeclaration{
|
|
PluginUniqueIdentifier: uniqueIdentifier.String(),
|
|
PluginID: uniqueIdentifier.PluginID(),
|
|
Declaration: declaration,
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
} else if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &declaration, nil
|
|
}
|
|
|
|
func (p *PluginManager) GetPackage(
|
|
plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
|
|
) ([]byte, error) {
|
|
file, err := p.packageBucket.Get(plugin_unique_identifier.String())
|
|
if err != nil {
|
|
if os.IsNotExist(err) {
|
|
return nil, errors.New("plugin package not found, please upload it firstly")
|
|
}
|
|
return nil, err
|
|
}
|
|
|
|
return file, nil
|
|
}
|
|
|
|
func (p *PluginManager) GetDeclaration(
|
|
plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
|
|
tenant_id string,
|
|
runtime_type plugin_entities.PluginRuntimeType,
|
|
) (
|
|
*plugin_entities.PluginDeclaration, error,
|
|
) {
|
|
return helper.CombinedGetPluginDeclaration(
|
|
plugin_unique_identifier, runtime_type,
|
|
)
|
|
}
|
|
|
|
var (
|
|
pluginAssetCache *lru.Cache[string, []byte]
|
|
)
|
|
|
|
func pluginAssetCacheKey(
|
|
pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
|
|
path string,
|
|
) string {
|
|
return fmt.Sprintf("%s/%s", pluginUniqueIdentifier.String(), path)
|
|
}
|
|
func (p *PluginManager) ExtractPluginAsset(
|
|
pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
|
|
path string,
|
|
) ([]byte, error) {
|
|
key := pluginAssetCacheKey(pluginUniqueIdentifier, path)
|
|
cached, ok := pluginAssetCache.Get(key)
|
|
if ok {
|
|
return cached, nil
|
|
}
|
|
pkgBytes, err := manager.GetPackage(pluginUniqueIdentifier)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
zipDecoder, err := decoder.NewZipPluginDecoder(pkgBytes)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
assets, err := zipDecoder.Assets()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
pluginAssetCache.Add(key, assets[path])
|
|
return assets[path], nil
|
|
}
|