server.go 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124
  1. package remote_manager
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "os/exec"
  7. "sync"
  8. "time"
  9. "github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager/media_manager"
  10. "github.com/langgenius/dify-plugin-daemon/internal/types/app"
  11. "github.com/langgenius/dify-plugin-daemon/internal/types/entities/plugin_entities"
  12. "github.com/langgenius/dify-plugin-daemon/internal/utils/stream"
  13. "github.com/panjf2000/gnet/v2"
  14. gnet_errors "github.com/panjf2000/gnet/v2/pkg/errors"
  15. )
  16. type RemotePluginServer struct {
  17. server *DifyServer
  18. }
  19. type RemotePluginServerInterface interface {
  20. Read() (plugin_entities.PluginFullDuplexLifetime, error)
  21. Next() bool
  22. Wrap(f func(plugin_entities.PluginFullDuplexLifetime))
  23. Stop() error
  24. Launch() error
  25. }
  26. // continue accepting new connections
  27. func (r *RemotePluginServer) Read() (plugin_entities.PluginFullDuplexLifetime, error) {
  28. if r.server.response == nil {
  29. return nil, errors.New("plugin server not started")
  30. }
  31. return r.server.response.Read()
  32. }
  33. // Next returns true if there are more connections to be read
  34. func (r *RemotePluginServer) Next() bool {
  35. if r.server.response == nil {
  36. return false
  37. }
  38. return r.server.response.Next()
  39. }
  40. // Wrap wraps the wrap method of stream response
  41. func (r *RemotePluginServer) Wrap(f func(plugin_entities.PluginFullDuplexLifetime)) {
  42. r.server.response.Async(f)
  43. }
  44. // Stop stops the server
  45. func (r *RemotePluginServer) Stop() error {
  46. if r.server.response == nil {
  47. return errors.New("plugin server not started")
  48. }
  49. r.server.response.Close()
  50. err := r.server.engine.Stop(context.Background())
  51. if err == gnet_errors.ErrEmptyEngine || err == gnet_errors.ErrEngineInShutdown {
  52. return nil
  53. }
  54. return err
  55. }
  56. // Launch starts the server
  57. func (r *RemotePluginServer) Launch() error {
  58. // kill the process if port is already in use
  59. exec.Command("fuser", "-k", "tcp", fmt.Sprintf("%d", r.server.port)).Run()
  60. time.Sleep(time.Millisecond * 100)
  61. err := gnet.Run(
  62. r.server, r.server.addr, gnet.WithMulticore(r.server.multicore),
  63. gnet.WithNumEventLoop(r.server.num_loops),
  64. )
  65. if err != nil {
  66. r.Stop()
  67. }
  68. return err
  69. }
  70. // NewRemotePluginServer creates a new RemotePluginServer
  71. func NewRemotePluginServer(config *app.Config, media_manager *media_manager.MediaManager) *RemotePluginServer {
  72. addr := fmt.Sprintf(
  73. "tcp://%s:%d",
  74. config.PluginRemoteInstallingHost,
  75. config.PluginRemoteInstallingPort,
  76. )
  77. response := stream.NewStream[plugin_entities.PluginFullDuplexLifetime](
  78. config.PluginRemoteInstallingMaxConn,
  79. )
  80. multicore := true
  81. s := &DifyServer{
  82. mediaManager: media_manager,
  83. addr: addr,
  84. port: config.PluginRemoteInstallingPort,
  85. multicore: multicore,
  86. num_loops: config.PluginRemoteInstallServerEventLoopNums,
  87. response: response,
  88. plugins: make(map[int]*RemotePluginRuntime),
  89. plugins_lock: &sync.RWMutex{},
  90. shutdown_chan: make(chan bool),
  91. max_conn: int32(config.PluginRemoteInstallingMaxConn),
  92. }
  93. manager := &RemotePluginServer{
  94. server: s,
  95. }
  96. return manager
  97. }