install_plugin.go 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494
  1. package service
  2. import (
  3. "fmt"
  4. "os"
  5. "github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager"
  6. "github.com/langgenius/dify-plugin-daemon/internal/core/plugin_packager/decoder"
  7. "github.com/langgenius/dify-plugin-daemon/internal/db"
  8. "github.com/langgenius/dify-plugin-daemon/internal/types/app"
  9. "github.com/langgenius/dify-plugin-daemon/internal/types/entities"
  10. "github.com/langgenius/dify-plugin-daemon/internal/types/entities/plugin_entities"
  11. "github.com/langgenius/dify-plugin-daemon/internal/types/models"
  12. "github.com/langgenius/dify-plugin-daemon/internal/types/models/curd"
  13. "github.com/langgenius/dify-plugin-daemon/internal/utils/log"
  14. "github.com/langgenius/dify-plugin-daemon/internal/utils/routine"
  15. "github.com/langgenius/dify-plugin-daemon/internal/utils/stream"
  16. "gorm.io/gorm"
  17. )
  18. type InstallPluginResponse struct {
  19. AllInstalled bool `json:"all_installed"`
  20. TaskID string `json:"task_id"`
  21. }
  22. type InstallPluginOnDoneHandler func(
  23. plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  24. declaration *plugin_entities.PluginDeclaration,
  25. ) error
  26. func InstallPluginRuntimeToTenant(
  27. config *app.Config,
  28. tenant_id string,
  29. plugin_unique_identifiers []plugin_entities.PluginUniqueIdentifier,
  30. source string,
  31. meta map[string]any,
  32. on_done InstallPluginOnDoneHandler, // since installing plugin is a async task, we need to call it asynchronously
  33. ) (*InstallPluginResponse, error) {
  34. response := &InstallPluginResponse{}
  35. plugins_wait_for_installation := []plugin_entities.PluginUniqueIdentifier{}
  36. task := &models.InstallTask{
  37. Status: models.InstallTaskStatusRunning,
  38. TenantID: tenant_id,
  39. TotalPlugins: len(plugin_unique_identifiers),
  40. CompletedPlugins: 0,
  41. Plugins: []models.InstallTaskPluginStatus{},
  42. }
  43. for i, plugin_unique_identifier := range plugin_unique_identifiers {
  44. // check if plugin is already installed
  45. plugin, err := db.GetOne[models.Plugin](
  46. db.Equal("plugin_unique_identifier", plugin_unique_identifier.String()),
  47. )
  48. task.Plugins = append(task.Plugins, models.InstallTaskPluginStatus{
  49. PluginUniqueIdentifier: plugin_unique_identifier,
  50. PluginID: plugin_unique_identifier.PluginID(),
  51. Status: models.InstallTaskStatusPending,
  52. Message: "",
  53. })
  54. if err == nil {
  55. // already installed by other tenant
  56. declaration := plugin.Declaration
  57. if _, _, err := curd.InstallPlugin(
  58. tenant_id,
  59. plugin_unique_identifier,
  60. plugin.InstallType,
  61. &declaration,
  62. source,
  63. meta,
  64. ); err != nil {
  65. return nil, err
  66. }
  67. task.CompletedPlugins++
  68. task.Plugins[i].Status = models.InstallTaskStatusSuccess
  69. task.Plugins[i].Message = "Installed"
  70. continue
  71. }
  72. if err != db.ErrDatabaseNotFound {
  73. return nil, err
  74. }
  75. plugins_wait_for_installation = append(plugins_wait_for_installation, plugin_unique_identifier)
  76. }
  77. if len(plugins_wait_for_installation) == 0 {
  78. response.AllInstalled = true
  79. response.TaskID = ""
  80. return response, nil
  81. }
  82. err := db.Create(task)
  83. if err != nil {
  84. return nil, err
  85. }
  86. response.TaskID = task.ID
  87. manager := plugin_manager.Manager()
  88. tasks := []func(){}
  89. for _, plugin_unique_identifier := range plugins_wait_for_installation {
  90. // copy the variable to avoid race condition
  91. plugin_unique_identifier := plugin_unique_identifier
  92. declaration, err := manager.GetDeclaration(plugin_unique_identifier)
  93. if err != nil {
  94. return nil, err
  95. }
  96. tasks = append(tasks, func() {
  97. updateTaskStatus := func(modifier func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus)) {
  98. if err := db.WithTransaction(func(tx *gorm.DB) error {
  99. task, err := db.GetOne[models.InstallTask](
  100. db.WithTransactionContext(tx),
  101. db.Equal("id", task.ID),
  102. db.WLock(), // write lock, multiple tasks can't update the same task
  103. )
  104. if err == db.ErrDatabaseNotFound {
  105. return nil
  106. }
  107. if err != nil {
  108. return err
  109. }
  110. task_pointer := &task
  111. var plugin_status *models.InstallTaskPluginStatus
  112. for i := range task.Plugins {
  113. if task.Plugins[i].PluginUniqueIdentifier == plugin_unique_identifier {
  114. plugin_status = &task.Plugins[i]
  115. break
  116. }
  117. }
  118. if plugin_status == nil {
  119. return nil
  120. }
  121. modifier(task_pointer, plugin_status)
  122. return db.Update(task_pointer, tx)
  123. }); err != nil {
  124. log.Error("failed to update install task status %s", err.Error())
  125. }
  126. }
  127. pkg_path, err := manager.GetPackagePath(plugin_unique_identifier)
  128. if err != nil {
  129. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  130. task.Status = models.InstallTaskStatusFailed
  131. plugin.Status = models.InstallTaskStatusFailed
  132. plugin.Message = err.Error()
  133. })
  134. return
  135. }
  136. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  137. plugin.Status = models.InstallTaskStatusRunning
  138. plugin.Message = "Installing"
  139. })
  140. var stream *stream.Stream[plugin_manager.PluginInstallResponse]
  141. if config.Platform == app.PLATFORM_AWS_LAMBDA {
  142. var zip_decoder *decoder.ZipPluginDecoder
  143. var pkg_file []byte
  144. pkg_file, err = os.ReadFile(pkg_path)
  145. if err != nil {
  146. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  147. task.Status = models.InstallTaskStatusFailed
  148. plugin.Status = models.InstallTaskStatusFailed
  149. plugin.Message = "Failed to read plugin package"
  150. })
  151. return
  152. }
  153. zip_decoder, err = decoder.NewZipPluginDecoder(pkg_file)
  154. if err != nil {
  155. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  156. task.Status = models.InstallTaskStatusFailed
  157. plugin.Status = models.InstallTaskStatusFailed
  158. plugin.Message = err.Error()
  159. })
  160. return
  161. }
  162. stream, err = manager.InstallToAWSFromPkg(zip_decoder, source, meta)
  163. } else if config.Platform == app.PLATFORM_LOCAL {
  164. stream, err = manager.InstallToLocal(pkg_path, plugin_unique_identifier, source, meta)
  165. } else {
  166. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  167. task.Status = models.InstallTaskStatusFailed
  168. plugin.Status = models.InstallTaskStatusFailed
  169. plugin.Message = "Unsupported platform"
  170. })
  171. return
  172. }
  173. if err != nil {
  174. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  175. task.Status = models.InstallTaskStatusFailed
  176. plugin.Status = models.InstallTaskStatusFailed
  177. plugin.Message = err.Error()
  178. })
  179. return
  180. }
  181. for stream.Next() {
  182. message, err := stream.Read()
  183. if err != nil {
  184. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  185. task.Status = models.InstallTaskStatusFailed
  186. plugin.Status = models.InstallTaskStatusFailed
  187. plugin.Message = err.Error()
  188. })
  189. return
  190. }
  191. if message.Event == plugin_manager.PluginInstallEventError {
  192. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  193. task.Status = models.InstallTaskStatusFailed
  194. plugin.Status = models.InstallTaskStatusFailed
  195. plugin.Message = message.Data
  196. })
  197. return
  198. }
  199. if message.Event == plugin_manager.PluginInstallEventDone {
  200. if err := on_done(plugin_unique_identifier, declaration); err != nil {
  201. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  202. task.Status = models.InstallTaskStatusFailed
  203. plugin.Status = models.InstallTaskStatusFailed
  204. plugin.Message = "Failed to create plugin"
  205. })
  206. return
  207. }
  208. }
  209. }
  210. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  211. plugin.Status = models.InstallTaskStatusSuccess
  212. plugin.Message = "Installed"
  213. task.CompletedPlugins++
  214. // check if all plugins are installed
  215. if task.CompletedPlugins == task.TotalPlugins {
  216. task.Status = models.InstallTaskStatusSuccess
  217. }
  218. })
  219. })
  220. }
  221. // submit async tasks
  222. routine.WithMaxRoutine(3, tasks)
  223. return response, nil
  224. }
  225. func InstallPluginFromIdentifiers(
  226. config *app.Config,
  227. tenant_id string,
  228. plugin_unique_identifiers []plugin_entities.PluginUniqueIdentifier,
  229. source string,
  230. meta map[string]any,
  231. ) *entities.Response {
  232. response, err := InstallPluginRuntimeToTenant(config, tenant_id, plugin_unique_identifiers, source, meta, func(
  233. plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  234. declaration *plugin_entities.PluginDeclaration,
  235. ) error {
  236. _, _, err := curd.InstallPlugin(
  237. tenant_id,
  238. plugin_unique_identifier,
  239. plugin_entities.PLUGIN_RUNTIME_TYPE_AWS,
  240. declaration,
  241. source,
  242. meta,
  243. )
  244. return err
  245. })
  246. if err != nil {
  247. return entities.NewErrorResponse(-500, err.Error())
  248. }
  249. return entities.NewSuccessResponse(response)
  250. }
  251. func UpgradePlugin(
  252. config *app.Config,
  253. tenant_id string,
  254. source string,
  255. meta map[string]any,
  256. original_plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  257. new_plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  258. ) *entities.Response {
  259. if original_plugin_unique_identifier == new_plugin_unique_identifier {
  260. return entities.NewErrorResponse(-400, "original and new plugin unique identifier are the same")
  261. }
  262. if original_plugin_unique_identifier.PluginID() != new_plugin_unique_identifier.PluginID() {
  263. return entities.NewErrorResponse(-400, "original and new plugin id are different")
  264. }
  265. // uninstall the original plugin
  266. installation, err := db.GetOne[models.PluginInstallation](
  267. db.Equal("tenant_id", tenant_id),
  268. db.Equal("plugin_unique_identifier", original_plugin_unique_identifier.String()),
  269. db.Equal("source", source),
  270. )
  271. if err == db.ErrDatabaseNotFound {
  272. return entities.NewErrorResponse(-404, "Plugin installation not found for this tenant")
  273. }
  274. if err != nil {
  275. return entities.NewErrorResponse(-500, err.Error())
  276. }
  277. // install the new plugin runtime
  278. response, err := InstallPluginRuntimeToTenant(
  279. config,
  280. tenant_id,
  281. []plugin_entities.PluginUniqueIdentifier{new_plugin_unique_identifier},
  282. source,
  283. meta,
  284. func(
  285. plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  286. declaration *plugin_entities.PluginDeclaration,
  287. ) error {
  288. // uninstall the original plugin
  289. err = curd.UpgradePlugin(
  290. tenant_id,
  291. original_plugin_unique_identifier,
  292. new_plugin_unique_identifier,
  293. declaration,
  294. plugin_entities.PluginRuntimeType(installation.RuntimeType),
  295. source,
  296. meta,
  297. )
  298. if err != nil {
  299. return err
  300. }
  301. return nil
  302. },
  303. )
  304. if err != nil {
  305. return entities.NewErrorResponse(-500, err.Error())
  306. }
  307. return entities.NewSuccessResponse(response)
  308. }
  309. func FetchPluginInstallationTasks(
  310. tenant_id string,
  311. page int,
  312. page_size int,
  313. ) *entities.Response {
  314. tasks, err := db.GetAll[models.InstallTask](
  315. db.Equal("tenant_id", tenant_id),
  316. db.OrderBy("created_at", true),
  317. db.Page(page, page_size),
  318. )
  319. if err != nil {
  320. return entities.NewErrorResponse(-500, err.Error())
  321. }
  322. return entities.NewSuccessResponse(tasks)
  323. }
  324. func FetchPluginInstallationTask(
  325. tenant_id string,
  326. task_id string,
  327. ) *entities.Response {
  328. task, err := db.GetOne[models.InstallTask](
  329. db.Equal("id", task_id),
  330. db.Equal("tenant_id", tenant_id),
  331. )
  332. if err != nil {
  333. return entities.NewErrorResponse(-500, err.Error())
  334. }
  335. return entities.NewSuccessResponse(task)
  336. }
  337. func DeletePluginInstallationTask(
  338. tenant_id string,
  339. task_id string,
  340. ) *entities.Response {
  341. err := db.DeleteByCondition(
  342. models.InstallTask{
  343. Model: models.Model{
  344. ID: task_id,
  345. },
  346. TenantID: tenant_id,
  347. },
  348. )
  349. if err != nil {
  350. return entities.NewErrorResponse(-500, err.Error())
  351. }
  352. return entities.NewSuccessResponse(true)
  353. }
  354. func DeletePluginInstallationItemFromTask(
  355. tenant_id string,
  356. task_id string,
  357. identifier plugin_entities.PluginUniqueIdentifier,
  358. ) *entities.Response {
  359. item, err := db.GetOne[models.InstallTask](
  360. db.Equal("task_id", task_id),
  361. db.Equal("tenant_id", tenant_id),
  362. )
  363. if err != nil {
  364. return entities.NewErrorResponse(-500, err.Error())
  365. }
  366. plugins := []models.InstallTaskPluginStatus{}
  367. for _, plugin := range item.Plugins {
  368. if plugin.PluginUniqueIdentifier != identifier {
  369. plugins = append(plugins, plugin)
  370. }
  371. }
  372. if len(plugins) == 0 {
  373. err = db.Delete(&item)
  374. } else {
  375. item.Plugins = plugins
  376. err = db.Update(&item)
  377. }
  378. if err != nil {
  379. return entities.NewErrorResponse(-500, err.Error())
  380. }
  381. return entities.NewSuccessResponse(true)
  382. }
  383. func FetchPluginFromIdentifier(
  384. plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  385. ) *entities.Response {
  386. _, err := db.GetOne[models.Plugin](
  387. db.Equal("plugin_unique_identifier", plugin_unique_identifier.String()),
  388. )
  389. if err == db.ErrDatabaseNotFound {
  390. return entities.NewSuccessResponse(false)
  391. }
  392. if err != nil {
  393. return entities.NewErrorResponse(-500, err.Error())
  394. }
  395. return entities.NewSuccessResponse(true)
  396. }
  397. func UninstallPlugin(
  398. tenant_id string,
  399. plugin_installation_id string,
  400. ) *entities.Response {
  401. // Check if the plugin exists for the tenant
  402. installation, err := db.GetOne[models.PluginInstallation](
  403. db.Equal("tenant_id", tenant_id),
  404. db.Equal("id", plugin_installation_id),
  405. )
  406. if err == db.ErrDatabaseNotFound {
  407. return entities.NewErrorResponse(-404, "Plugin installation not found for this tenant")
  408. }
  409. if err != nil {
  410. return entities.NewErrorResponse(-500, err.Error())
  411. }
  412. plugin_unique_identifier, err := plugin_entities.NewPluginUniqueIdentifier(installation.PluginUniqueIdentifier)
  413. if err != nil {
  414. return entities.NewErrorResponse(-500, fmt.Sprintf("failed to parse plugin unique identifier: %v", err))
  415. }
  416. // Uninstall the plugin
  417. _, err = curd.UninstallPlugin(
  418. tenant_id,
  419. plugin_unique_identifier,
  420. installation.ID,
  421. )
  422. if err != nil {
  423. return entities.NewErrorResponse(-500, fmt.Sprintf("Failed to uninstall plugin: %s", err.Error()))
  424. }
  425. return entities.NewSuccessResponse(true)
  426. }