mirror of
https://github.com/langgenius/dify-plugin-daemon.git
synced 2026-07-25 13:35:32 -04:00
878edde455
* 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>
198 lines
7.2 KiB
Go
198 lines
7.2 KiB
Go
package server
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net/http"
|
|
"time"
|
|
|
|
"github.com/gin-gonic/gin"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/plugin_daemon/backwards_invocation/transaction"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/server/controllers"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/service"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/types/app"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/utils/log"
|
|
|
|
sentrygin "github.com/getsentry/sentry-go/gin"
|
|
)
|
|
|
|
// server starts a http server and returns a function to stop it
|
|
func (app *App) server(config *app.Config) func() {
|
|
engine := gin.New()
|
|
if *config.HealthApiLogEnabled {
|
|
engine.Use(gin.Logger())
|
|
} else {
|
|
engine.Use(gin.LoggerWithConfig(gin.LoggerConfig{
|
|
SkipPaths: []string{"/health/check"},
|
|
}))
|
|
}
|
|
engine.Use(gin.Recovery())
|
|
engine.Use(controllers.CollectActiveRequests())
|
|
engine.GET("/health/check", controllers.HealthCheck(config))
|
|
|
|
endpointGroup := engine.Group("/e")
|
|
serverlessTransactionGroup := engine.Group("/backwards-invocation")
|
|
pluginGroup := engine.Group("/plugin/:tenant_id")
|
|
pprofGroup := engine.Group("/debug/pprof")
|
|
|
|
if config.AdminApiEnabled {
|
|
if len(config.AdminApiKey) < 10 {
|
|
log.Panic("length of admin api key must be greater than 10")
|
|
}
|
|
|
|
adminGroup := engine.Group("/admin")
|
|
adminGroup.Use(app.AdminAPIKey(config.AdminApiKey))
|
|
|
|
app.adminGroup(adminGroup, config)
|
|
}
|
|
|
|
if config.SentryEnabled {
|
|
// setup sentry for all groups
|
|
sentryGroup := []*gin.RouterGroup{
|
|
endpointGroup,
|
|
serverlessTransactionGroup,
|
|
pluginGroup,
|
|
}
|
|
for _, group := range sentryGroup {
|
|
group.Use(sentrygin.New(sentrygin.Options{
|
|
Repanic: true,
|
|
}))
|
|
}
|
|
}
|
|
|
|
app.endpointGroup(endpointGroup, config)
|
|
app.serverlessTransactionGroup(serverlessTransactionGroup, config)
|
|
app.pluginGroup(pluginGroup, config)
|
|
app.pprofGroup(pprofGroup, config)
|
|
|
|
srv := &http.Server{
|
|
Addr: fmt.Sprintf(":%d", config.ServerPort),
|
|
Handler: engine,
|
|
}
|
|
|
|
go func() {
|
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
log.Panic("listen: %s\n", err)
|
|
}
|
|
}()
|
|
|
|
return func() {
|
|
if err := srv.Shutdown(context.Background()); err != nil {
|
|
log.Panic("Server Shutdown: %s\n", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (app *App) pluginGroup(group *gin.RouterGroup, config *app.Config) {
|
|
group.Use(CheckingKey(config.ServerKey))
|
|
|
|
app.remoteDebuggingGroup(group.Group("/debugging"), config)
|
|
app.pluginDispatchGroup(group.Group("/dispatch"), config)
|
|
app.pluginManagementGroup(group.Group("/management"), config)
|
|
app.endpointManagementGroup(group.Group("/endpoint"))
|
|
app.pluginAssetGroup(group.Group("/asset"))
|
|
}
|
|
|
|
func (app *App) pluginDispatchGroup(group *gin.RouterGroup, config *app.Config) {
|
|
group.Use(controllers.CollectActiveDispatchRequests())
|
|
group.Use(app.FetchPluginInstallation())
|
|
group.Use(app.RedirectPluginInvoke())
|
|
group.Use(app.InitClusterID())
|
|
|
|
group.POST("/agent_strategy/invoke", controllers.InvokeAgentStrategy(config))
|
|
|
|
app.setupGeneratedRoutes(group, config)
|
|
}
|
|
|
|
func (app *App) remoteDebuggingGroup(group *gin.RouterGroup, config *app.Config) {
|
|
if config.PluginRemoteInstallingEnabled != nil && *config.PluginRemoteInstallingEnabled {
|
|
group.POST("/key", CheckingKey(config.ServerKey), controllers.GetRemoteDebuggingKey)
|
|
}
|
|
}
|
|
|
|
func (app *App) endpointGroup(group *gin.RouterGroup, config *app.Config) {
|
|
if config.PluginEndpointEnabled != nil && *config.PluginEndpointEnabled {
|
|
group.HEAD("/:hook_id/*path", app.Endpoint(config))
|
|
group.POST("/:hook_id/*path", app.Endpoint(config))
|
|
group.GET("/:hook_id/*path", app.Endpoint(config))
|
|
group.PUT("/:hook_id/*path", app.Endpoint(config))
|
|
group.DELETE("/:hook_id/*path", app.Endpoint(config))
|
|
group.OPTIONS("/:hook_id/*path", app.Endpoint(config))
|
|
}
|
|
}
|
|
|
|
func (appRef *App) serverlessTransactionGroup(group *gin.RouterGroup, config *app.Config) {
|
|
if config.Platform == app.PLATFORM_SERVERLESS {
|
|
appRef.serverlessTransactionHandler = transaction.NewServerlessTransactionHandler(
|
|
time.Duration(config.MaxServerlessTransactionTimeout) * time.Second,
|
|
)
|
|
group.POST(
|
|
"/transaction",
|
|
service.HandleServerlessPluginTransaction(appRef.serverlessTransactionHandler),
|
|
)
|
|
}
|
|
}
|
|
|
|
func (app *App) endpointManagementGroup(group *gin.RouterGroup) {
|
|
group.POST("/setup", controllers.SetupEndpoint)
|
|
group.POST("/remove", controllers.RemoveEndpoint)
|
|
group.POST("/update", controllers.UpdateEndpoint)
|
|
group.GET("/list", controllers.ListEndpoints)
|
|
group.GET("/list/plugin", controllers.ListPluginEndpoints)
|
|
group.POST("/enable", controllers.EnableEndpoint)
|
|
group.POST("/disable", controllers.DisableEndpoint)
|
|
}
|
|
|
|
func (app *App) pluginManagementGroup(group *gin.RouterGroup, config *app.Config) {
|
|
group.POST("/install/upload/package", controllers.UploadPlugin(config))
|
|
group.POST("/install/upload/bundle", controllers.UploadBundle(config))
|
|
group.POST("/install/identifiers", controllers.InstallPluginFromIdentifiers(config))
|
|
group.POST("/install/upgrade", controllers.UpgradePlugin(config))
|
|
group.GET("/install/tasks/:id", controllers.FetchPluginInstallationTask)
|
|
group.POST("/install/tasks/delete_all", controllers.DeleteAllPluginInstallationTasks)
|
|
group.POST("/install/tasks/:id/delete", controllers.DeletePluginInstallationTask)
|
|
group.POST("/install/tasks/:id/delete/*identifier", controllers.DeletePluginInstallationItemFromTask)
|
|
group.GET("/install/tasks", controllers.FetchPluginInstallationTasks)
|
|
group.GET("/decode/from_identifier", controllers.DecodePluginFromIdentifier(config))
|
|
group.GET("/fetch/manifest", controllers.FetchPluginManifest)
|
|
group.GET("/fetch/identifier", controllers.FetchPluginFromIdentifier)
|
|
group.POST("/uninstall", controllers.UninstallPlugin)
|
|
group.GET("/list", controllers.ListPlugins)
|
|
group.POST("/installation/fetch/batch", controllers.BatchFetchPluginInstallationByIDs)
|
|
group.POST("/installation/missing", controllers.FetchMissingPluginInstallations)
|
|
group.GET("/models", controllers.ListModels)
|
|
group.GET("/tools", controllers.ListTools)
|
|
group.GET("/tool", controllers.GetTool)
|
|
group.POST("/tools/check_existence", controllers.CheckToolExistence)
|
|
group.GET("/agent_strategies", controllers.ListAgentStrategies)
|
|
group.GET("/agent_strategy", controllers.GetAgentStrategy)
|
|
group.GET("/datasources", controllers.ListDatasources)
|
|
group.GET("/datasource", controllers.GetDatasource)
|
|
}
|
|
|
|
func (app *App) adminGroup(group *gin.RouterGroup, config *app.Config) {
|
|
group.POST("/plugin/serverless/reinstall", controllers.ReinstallPluginFromIdentifier(config))
|
|
}
|
|
|
|
func (app *App) pluginAssetGroup(group *gin.RouterGroup) {
|
|
group.GET("/:id", controllers.GetAsset)
|
|
}
|
|
|
|
func (app *App) pprofGroup(group *gin.RouterGroup, config *app.Config) {
|
|
if config.PPROFEnabled {
|
|
group.Use(CheckingKey(config.ServerKey))
|
|
|
|
group.GET("/", controllers.PprofIndex)
|
|
group.GET("/cmdline", controllers.PprofCmdline)
|
|
group.GET("/profile", controllers.PprofProfile)
|
|
group.GET("/symbol", controllers.PprofSymbol)
|
|
group.GET("/trace", controllers.PprofTrace)
|
|
group.GET("/goroutine", controllers.PprofGoroutine)
|
|
group.GET("/heap", controllers.PprofHeap)
|
|
group.GET("/allocs", controllers.PprofAllocs)
|
|
group.GET("/block", controllers.PprofBlock)
|
|
group.GET("/mutex", controllers.PprofMutex)
|
|
group.GET("/threadcreate", controllers.PprofThreadcreate)
|
|
}
|
|
}
|