mirror of
https://github.com/langgenius/dify-plugin-daemon.git
synced 2026-07-22 01:35:24 -04:00
4314b8f1ee
Remote debugging plugins were not being synchronized across cluster nodes,
causing "no plugin available nodes found" errors when trying to invoke
plugins from different nodes.
1. **Remote debugging plugins not registered to cluster** - The
`ClusterTunnel` notifier was not being added to ControlPanel
2. **Plugin ID inconsistency** - Remote plugins used different plugin_id
formats during installation vs. querying
3. **Non-idempotent registration** - `RegisterPlugin` failed on reconnection
with "plugin has been registered" error
- **internal/types/models/curd/atomic.go**:
- Unify plugin_id calculation for remote plugins (author/name without version)
- Remove plugin_id from plugin query conditions
- Clear old cache when plugin_id is updated
- **internal/cluster/plugin.go**:
- Make `RegisterPlugin` idempotent by updating existing plugin instead
of returning error
- **internal/core/control_panel/daemon.go**:
- Add cluster field to ControlPanel
- Add SetCluster() method for lazy cluster initialization
- **internal/core/control_panel/server_debugger.go**:
- Register remote debugging plugins to cluster on connection
- Unregister from cluster on disconnection
- **internal/core/plugin_manager/manager.go**:
- Add SetCluster() method to set cluster after initialization
- **internal/server/server.go**:
- Call SetCluster() instead of AddClusterTunnel()
Only remote debugging plugins are synchronized across cluster nodes.
Local plugins run only on the node where they are installed and are
not registered to the cluster.
- Error handling improvements using `errors.Is()` instead of `==`
- Handle 404 for missing plugin assets gracefully
- Handle already-installed debugging plugins gracefully
- Remote debugging plugin can be invoked from any node in the cluster
- Plugin reconnection works without errors
- Cache invalidation works correctly when plugin_id changes
243 lines
5.8 KiB
Go
243 lines
5.8 KiB
Go
package cluster
|
|
|
|
import (
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/io_tunnel/access_types"
|
|
"github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager/basic_runtime"
|
|
"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"
|
|
)
|
|
|
|
type fakePlugin struct {
|
|
plugin_entities.PluginRuntime
|
|
basic_runtime.BasicChecksum
|
|
}
|
|
|
|
func (r *fakePlugin) InitEnvironment() error {
|
|
return nil
|
|
}
|
|
|
|
func (r *fakePlugin) Checksum() (string, error) {
|
|
return "", nil
|
|
}
|
|
|
|
func (r *fakePlugin) Identity() (plugin_entities.PluginUniqueIdentifier, error) {
|
|
return plugin_entities.PluginUniqueIdentifier(""), nil
|
|
}
|
|
|
|
func (r *fakePlugin) StartPlugin() error {
|
|
return nil
|
|
}
|
|
|
|
func (r *fakePlugin) Type() plugin_entities.PluginRuntimeType {
|
|
return plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
|
|
}
|
|
|
|
func (r *fakePlugin) Wait() (<-chan bool, error) {
|
|
return nil, nil
|
|
}
|
|
|
|
func (r *fakePlugin) Listen(string) (*entities.Broadcast[plugin_entities.SessionMessage], error) {
|
|
return nil, nil
|
|
}
|
|
|
|
func (r *fakePlugin) Write(string, access_types.PluginAccessAction, []byte) error {
|
|
return nil
|
|
}
|
|
|
|
func getRandomPluginRuntime() fakePlugin {
|
|
return fakePlugin{
|
|
PluginRuntime: plugin_entities.PluginRuntime{
|
|
Config: plugin_entities.PluginDeclaration{
|
|
PluginDeclarationWithoutAdvancedFields: plugin_entities.PluginDeclarationWithoutAdvancedFields{
|
|
Name: uuid.New().String(),
|
|
Label: plugin_entities.I18nObject{
|
|
EnUS: "label",
|
|
},
|
|
Version: "0.0.1",
|
|
Type: manifest_entities.PluginType,
|
|
Author: "Yeuoly",
|
|
CreatedAt: time.Now(),
|
|
Plugins: plugin_entities.PluginExtensions{
|
|
Tools: []string{"test"},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
}
|
|
}
|
|
|
|
func TestPluginScheduleLifetime(t *testing.T) {
|
|
plugin := getRandomPluginRuntime()
|
|
cluster, err := createSimulationCluster(1)
|
|
if err != nil {
|
|
t.Errorf("create simulation cluster failed: %v", err)
|
|
return
|
|
}
|
|
|
|
launchSimulationCluster(cluster)
|
|
defer closeSimulationCluster(cluster, t)
|
|
|
|
time.Sleep(time.Second * 1)
|
|
|
|
// add plugin to the cluster
|
|
err = cluster[0].RegisterPlugin(&plugin)
|
|
if err != nil {
|
|
t.Errorf("register plugin failed: %v", err)
|
|
return
|
|
}
|
|
|
|
identity, err := plugin.Identity()
|
|
if err != nil {
|
|
t.Errorf("get plugin identity failed: %v", err)
|
|
return
|
|
}
|
|
|
|
hashedIdentity := plugin_entities.HashedIdentity(identity.String())
|
|
|
|
nodes, err := cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 1 {
|
|
t.Errorf("plugin not scheduled")
|
|
return
|
|
}
|
|
|
|
if nodes[0] != cluster[0].id {
|
|
t.Errorf("plugin scheduled to wrong node")
|
|
return
|
|
}
|
|
|
|
// trigger plugin stop
|
|
plugin.Stop()
|
|
|
|
// notify plugin has been stopped
|
|
if err := cluster[0].UnregisterPlugin(&plugin); err != nil {
|
|
t.Errorf("unregister plugin failed: %v", err)
|
|
return
|
|
}
|
|
|
|
// wait for the plugin to stop
|
|
time.Sleep(time.Second * 1)
|
|
|
|
// check if the plugin is stopped
|
|
nodes, err = cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 0 {
|
|
t.Errorf("plugin not stopped")
|
|
return
|
|
}
|
|
}
|
|
|
|
func TestPluginRegisterIdempotent(t *testing.T) {
|
|
plugin := getRandomPluginRuntime()
|
|
cluster, err := createSimulationCluster(1)
|
|
if err != nil {
|
|
t.Errorf("create simulation cluster failed: %v", err)
|
|
return
|
|
}
|
|
|
|
launchSimulationCluster(cluster)
|
|
defer closeSimulationCluster(cluster, t)
|
|
|
|
// wait for cluster to be ready
|
|
time.Sleep(time.Second * 1)
|
|
|
|
// first registration should succeed
|
|
err = cluster[0].RegisterPlugin(&plugin)
|
|
if err != nil {
|
|
t.Errorf("first register plugin failed: %v", err)
|
|
return
|
|
}
|
|
|
|
identity, err := plugin.Identity()
|
|
if err != nil {
|
|
t.Errorf("get plugin identity failed: %v", err)
|
|
return
|
|
}
|
|
|
|
hashedIdentity := plugin_entities.HashedIdentity(identity.String())
|
|
|
|
// wait for plugin to be scheduled
|
|
time.Sleep(time.Second * 1)
|
|
|
|
// verify plugin is registered
|
|
nodes, err := cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 1 {
|
|
t.Errorf("plugin not scheduled after first registration")
|
|
return
|
|
}
|
|
|
|
// second registration with same identity should be idempotent (no error)
|
|
err = cluster[0].RegisterPlugin(&plugin)
|
|
if err != nil {
|
|
t.Errorf("second register plugin failed (should be idempotent): %v", err)
|
|
return
|
|
}
|
|
|
|
// verify plugin is still registered after second registration
|
|
nodes, err = cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed after second registration: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 1 {
|
|
t.Errorf("plugin not available after second registration")
|
|
return
|
|
}
|
|
|
|
// unregister the plugin
|
|
if err := cluster[0].UnregisterPlugin(&plugin); err != nil {
|
|
t.Errorf("unregister plugin failed: %v", err)
|
|
return
|
|
}
|
|
|
|
// verify plugin is unregistered
|
|
nodes, err = cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed after unregister: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 0 {
|
|
t.Errorf("plugin still available after unregister")
|
|
return
|
|
}
|
|
|
|
// registration after unregister should succeed again
|
|
err = cluster[0].RegisterPlugin(&plugin)
|
|
if err != nil {
|
|
t.Errorf("register after unregister failed: %v", err)
|
|
return
|
|
}
|
|
|
|
// verify plugin is registered again
|
|
nodes, err = cluster[0].FetchPluginAvailableNodesByHashedId(hashedIdentity)
|
|
if err != nil {
|
|
t.Errorf("fetch plugin available nodes failed after re-registration: %v", err)
|
|
return
|
|
}
|
|
|
|
if len(nodes) != 1 {
|
|
t.Errorf("plugin not available after re-registration")
|
|
return
|
|
}
|
|
}
|