server.go 3.2 KB

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