install_plugin.go 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627
  1. package service
  2. import (
  3. "errors"
  4. "fmt"
  5. "time"
  6. "github.com/langgenius/dify-plugin-daemon/internal/core/plugin_manager"
  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/exception"
  10. "github.com/langgenius/dify-plugin-daemon/internal/types/models"
  11. "github.com/langgenius/dify-plugin-daemon/internal/types/models/curd"
  12. "github.com/langgenius/dify-plugin-daemon/internal/utils/cache/helper"
  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. "github.com/langgenius/dify-plugin-daemon/pkg/entities"
  17. "github.com/langgenius/dify-plugin-daemon/pkg/entities/plugin_entities"
  18. "github.com/langgenius/dify-plugin-daemon/pkg/plugin_packager/decoder"
  19. "gorm.io/gorm"
  20. )
  21. type InstallPluginResponse struct {
  22. AllInstalled bool `json:"all_installed"`
  23. TaskID string `json:"task_id"`
  24. }
  25. type InstallPluginOnDoneHandler func(
  26. pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
  27. declaration *plugin_entities.PluginDeclaration,
  28. meta map[string]any,
  29. ) error
  30. func InstallPluginRuntimeToTenant(
  31. config *app.Config,
  32. tenant_id string,
  33. plugin_unique_identifiers []plugin_entities.PluginUniqueIdentifier,
  34. source string,
  35. metas []map[string]any,
  36. onDone InstallPluginOnDoneHandler, // since installing plugin is a async task, we need to call it asynchronously
  37. ) (*InstallPluginResponse, error) {
  38. response := &InstallPluginResponse{}
  39. pluginsWaitForInstallation := []plugin_entities.PluginUniqueIdentifier{}
  40. runtimeType := plugin_entities.PluginRuntimeType("")
  41. if config.Platform == app.PLATFORM_SERVERLESS {
  42. runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_SERVERLESS
  43. } else if config.Platform == app.PLATFORM_LOCAL {
  44. runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
  45. } else {
  46. return nil, fmt.Errorf("unsupported platform: %s", config.Platform)
  47. }
  48. task := &models.InstallTask{
  49. Status: models.InstallTaskStatusRunning,
  50. TenantID: tenant_id,
  51. TotalPlugins: len(plugin_unique_identifiers),
  52. CompletedPlugins: 0,
  53. Plugins: []models.InstallTaskPluginStatus{},
  54. }
  55. for i, pluginUniqueIdentifier := range plugin_unique_identifiers {
  56. // fetch plugin declaration first, before installing, we need to ensure pkg is uploaded
  57. pluginDeclaration, err := helper.CombinedGetPluginDeclaration(
  58. pluginUniqueIdentifier,
  59. runtimeType,
  60. )
  61. if err != nil {
  62. return nil, err
  63. }
  64. // check if plugin is already installed
  65. _, err = db.GetOne[models.Plugin](
  66. db.Equal("plugin_unique_identifier", pluginUniqueIdentifier.String()),
  67. )
  68. task.Plugins = append(task.Plugins, models.InstallTaskPluginStatus{
  69. PluginUniqueIdentifier: pluginUniqueIdentifier,
  70. PluginID: pluginUniqueIdentifier.PluginID(),
  71. Status: models.InstallTaskStatusPending,
  72. Icon: pluginDeclaration.Icon,
  73. Labels: pluginDeclaration.Label,
  74. Message: "",
  75. })
  76. if err == nil {
  77. if err := onDone(pluginUniqueIdentifier, pluginDeclaration, metas[i]); err != nil {
  78. return nil, errors.Join(err, errors.New("failed on plugin installation"))
  79. } else {
  80. task.CompletedPlugins++
  81. task.Plugins[i].Status = models.InstallTaskStatusSuccess
  82. task.Plugins[i].Message = "Installed"
  83. }
  84. continue
  85. }
  86. if err != db.ErrDatabaseNotFound {
  87. return nil, err
  88. }
  89. pluginsWaitForInstallation = append(pluginsWaitForInstallation, pluginUniqueIdentifier)
  90. }
  91. if len(pluginsWaitForInstallation) == 0 {
  92. response.AllInstalled = true
  93. response.TaskID = ""
  94. return response, nil
  95. }
  96. err := db.Create(task)
  97. if err != nil {
  98. return nil, err
  99. }
  100. response.TaskID = task.ID
  101. manager := plugin_manager.Manager()
  102. tasks := []func(){}
  103. for i, pluginUniqueIdentifier := range pluginsWaitForInstallation {
  104. // copy the variable to avoid race condition
  105. pluginUniqueIdentifier := pluginUniqueIdentifier
  106. declaration, err := helper.CombinedGetPluginDeclaration(
  107. pluginUniqueIdentifier,
  108. runtimeType,
  109. )
  110. if err != nil {
  111. return nil, err
  112. }
  113. i := i
  114. tasks = append(tasks, func() {
  115. updateTaskStatus := func(modifier func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus)) {
  116. if err := db.WithTransaction(func(tx *gorm.DB) error {
  117. task, err := db.GetOne[models.InstallTask](
  118. db.WithTransactionContext(tx),
  119. db.Equal("id", task.ID),
  120. db.WLock(), // write lock, multiple tasks can't update the same task
  121. )
  122. if err == db.ErrDatabaseNotFound {
  123. return nil
  124. }
  125. if err != nil {
  126. return err
  127. }
  128. taskPointer := &task
  129. var pluginStatus *models.InstallTaskPluginStatus
  130. for i := range task.Plugins {
  131. if task.Plugins[i].PluginUniqueIdentifier == pluginUniqueIdentifier {
  132. pluginStatus = &task.Plugins[i]
  133. break
  134. }
  135. }
  136. if pluginStatus == nil {
  137. return nil
  138. }
  139. modifier(taskPointer, pluginStatus)
  140. successes := 0
  141. for _, plugin := range taskPointer.Plugins {
  142. if plugin.Status == models.InstallTaskStatusSuccess {
  143. successes++
  144. }
  145. }
  146. if successes == len(taskPointer.Plugins) {
  147. // update status
  148. taskPointer.Status = models.InstallTaskStatusSuccess
  149. // delete the task after 120 seconds without transaction
  150. time.AfterFunc(120*time.Second, func() {
  151. db.Delete(taskPointer)
  152. })
  153. }
  154. return db.Update(taskPointer, tx)
  155. }); err != nil {
  156. log.Error("failed to update install task status %s", err.Error())
  157. }
  158. }
  159. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  160. plugin.Status = models.InstallTaskStatusRunning
  161. plugin.Message = "Installing"
  162. })
  163. var stream *stream.Stream[plugin_manager.PluginInstallResponse]
  164. if config.Platform == app.PLATFORM_SERVERLESS {
  165. var zipDecoder *decoder.ZipPluginDecoder
  166. var pkgFile []byte
  167. pkgFile, err = manager.GetPackage(pluginUniqueIdentifier)
  168. if err != nil {
  169. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  170. task.Status = models.InstallTaskStatusFailed
  171. plugin.Status = models.InstallTaskStatusFailed
  172. plugin.Message = "Failed to read plugin package"
  173. })
  174. return
  175. }
  176. zipDecoder, err = decoder.NewZipPluginDecoder(pkgFile)
  177. if err != nil {
  178. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  179. task.Status = models.InstallTaskStatusFailed
  180. plugin.Status = models.InstallTaskStatusFailed
  181. plugin.Message = err.Error()
  182. })
  183. return
  184. }
  185. stream, err = manager.InstallToAWSFromPkg(pkgFile, zipDecoder, source, metas[i])
  186. } else if config.Platform == app.PLATFORM_LOCAL {
  187. stream, err = manager.InstallToLocal(pluginUniqueIdentifier, source, metas[i])
  188. } else {
  189. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  190. task.Status = models.InstallTaskStatusFailed
  191. plugin.Status = models.InstallTaskStatusFailed
  192. plugin.Message = "Unsupported platform"
  193. })
  194. return
  195. }
  196. if err != nil {
  197. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  198. task.Status = models.InstallTaskStatusFailed
  199. plugin.Status = models.InstallTaskStatusFailed
  200. plugin.Message = err.Error()
  201. })
  202. return
  203. }
  204. for stream.Next() {
  205. message, err := stream.Read()
  206. if err != nil {
  207. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  208. task.Status = models.InstallTaskStatusFailed
  209. plugin.Status = models.InstallTaskStatusFailed
  210. plugin.Message = err.Error()
  211. })
  212. return
  213. }
  214. if message.Event == plugin_manager.PluginInstallEventError {
  215. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  216. task.Status = models.InstallTaskStatusFailed
  217. plugin.Status = models.InstallTaskStatusFailed
  218. plugin.Message = message.Data
  219. })
  220. return
  221. }
  222. if message.Event == plugin_manager.PluginInstallEventDone {
  223. if err := onDone(pluginUniqueIdentifier, declaration, metas[i]); err != nil {
  224. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  225. task.Status = models.InstallTaskStatusFailed
  226. plugin.Status = models.InstallTaskStatusFailed
  227. plugin.Message = "Failed to create plugin, perhaps it's already installed"
  228. })
  229. return
  230. }
  231. }
  232. }
  233. updateTaskStatus(func(task *models.InstallTask, plugin *models.InstallTaskPluginStatus) {
  234. plugin.Status = models.InstallTaskStatusSuccess
  235. plugin.Message = "Installed"
  236. task.CompletedPlugins++
  237. // check if all plugins are installed
  238. if task.CompletedPlugins == task.TotalPlugins {
  239. task.Status = models.InstallTaskStatusSuccess
  240. }
  241. })
  242. })
  243. }
  244. // submit async tasks
  245. routine.WithMaxRoutine(5, tasks)
  246. return response, nil
  247. }
  248. func InstallPluginFromIdentifiers(
  249. config *app.Config,
  250. tenant_id string,
  251. plugin_unique_identifiers []plugin_entities.PluginUniqueIdentifier,
  252. source string,
  253. metas []map[string]any,
  254. ) *entities.Response {
  255. response, err := InstallPluginRuntimeToTenant(
  256. config,
  257. tenant_id,
  258. plugin_unique_identifiers,
  259. source,
  260. metas,
  261. func(
  262. pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
  263. declaration *plugin_entities.PluginDeclaration,
  264. meta map[string]any,
  265. ) error {
  266. runtimeType := plugin_entities.PluginRuntimeType("")
  267. switch config.Platform {
  268. case app.PLATFORM_SERVERLESS:
  269. runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_SERVERLESS
  270. case app.PLATFORM_LOCAL:
  271. runtimeType = plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL
  272. default:
  273. return fmt.Errorf("unsupported platform: %s", config.Platform)
  274. }
  275. _, _, err := curd.InstallPlugin(
  276. tenant_id,
  277. pluginUniqueIdentifier,
  278. runtimeType,
  279. declaration,
  280. source,
  281. meta,
  282. )
  283. return err
  284. },
  285. )
  286. if err != nil {
  287. if errors.Is(err, curd.ErrPluginAlreadyInstalled) {
  288. return exception.BadRequestError(err).ToResponse()
  289. }
  290. return exception.InternalServerError(err).ToResponse()
  291. }
  292. return entities.NewSuccessResponse(response)
  293. }
  294. func UpgradePlugin(
  295. config *app.Config,
  296. tenant_id string,
  297. source string,
  298. meta map[string]any,
  299. original_plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  300. new_plugin_unique_identifier plugin_entities.PluginUniqueIdentifier,
  301. ) *entities.Response {
  302. if original_plugin_unique_identifier == new_plugin_unique_identifier {
  303. return exception.BadRequestError(errors.New("original and new plugin unique identifier are the same")).ToResponse()
  304. }
  305. if original_plugin_unique_identifier.PluginID() != new_plugin_unique_identifier.PluginID() {
  306. return exception.BadRequestError(errors.New("original and new plugin id are different")).ToResponse()
  307. }
  308. // uninstall the original plugin
  309. installation, err := db.GetOne[models.PluginInstallation](
  310. db.Equal("tenant_id", tenant_id),
  311. db.Equal("plugin_unique_identifier", original_plugin_unique_identifier.String()),
  312. db.Equal("source", source),
  313. )
  314. if err == db.ErrDatabaseNotFound {
  315. return exception.NotFoundError(errors.New("plugin installation not found for this tenant")).ToResponse()
  316. }
  317. if err != nil {
  318. return exception.InternalServerError(err).ToResponse()
  319. }
  320. // install the new plugin runtime
  321. response, err := InstallPluginRuntimeToTenant(
  322. config,
  323. tenant_id,
  324. []plugin_entities.PluginUniqueIdentifier{new_plugin_unique_identifier},
  325. source,
  326. []map[string]any{meta},
  327. func(
  328. pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
  329. declaration *plugin_entities.PluginDeclaration,
  330. meta map[string]any,
  331. ) error {
  332. originalDeclaration, err := helper.CombinedGetPluginDeclaration(
  333. original_plugin_unique_identifier,
  334. plugin_entities.PluginRuntimeType(installation.RuntimeType),
  335. )
  336. if err != nil {
  337. return err
  338. }
  339. newDeclaration, err := helper.CombinedGetPluginDeclaration(
  340. new_plugin_unique_identifier,
  341. plugin_entities.PluginRuntimeType(installation.RuntimeType),
  342. )
  343. if err != nil {
  344. return err
  345. }
  346. // uninstall the original plugin
  347. upgradeResponse, err := curd.UpgradePlugin(
  348. tenant_id,
  349. original_plugin_unique_identifier,
  350. new_plugin_unique_identifier,
  351. originalDeclaration,
  352. newDeclaration,
  353. plugin_entities.PluginRuntimeType(installation.RuntimeType),
  354. source,
  355. meta,
  356. )
  357. if err != nil {
  358. return err
  359. }
  360. if upgradeResponse.IsOriginalPluginDeleted {
  361. // delete the plugin if no installation left
  362. manager := plugin_manager.Manager()
  363. if string(upgradeResponse.DeletedPlugin.InstallType) == string(
  364. plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL,
  365. ) {
  366. err = manager.UninstallFromLocal(
  367. plugin_entities.PluginUniqueIdentifier(upgradeResponse.DeletedPlugin.PluginUniqueIdentifier),
  368. )
  369. if err != nil {
  370. return err
  371. }
  372. }
  373. }
  374. return nil
  375. },
  376. )
  377. if err != nil {
  378. return exception.InternalServerError(err).ToResponse()
  379. }
  380. return entities.NewSuccessResponse(response)
  381. }
  382. func FetchPluginInstallationTasks(
  383. tenant_id string,
  384. page int,
  385. page_size int,
  386. ) *entities.Response {
  387. tasks, err := db.GetAll[models.InstallTask](
  388. db.Equal("tenant_id", tenant_id),
  389. db.OrderBy("created_at", true),
  390. db.Page(page, page_size),
  391. )
  392. if err != nil {
  393. return exception.InternalServerError(err).ToResponse()
  394. }
  395. return entities.NewSuccessResponse(tasks)
  396. }
  397. func FetchPluginInstallationTask(
  398. tenant_id string,
  399. task_id string,
  400. ) *entities.Response {
  401. task, err := db.GetOne[models.InstallTask](
  402. db.Equal("id", task_id),
  403. db.Equal("tenant_id", tenant_id),
  404. )
  405. if err != nil {
  406. return exception.InternalServerError(err).ToResponse()
  407. }
  408. return entities.NewSuccessResponse(task)
  409. }
  410. func DeletePluginInstallationTask(
  411. tenant_id string,
  412. task_id string,
  413. ) *entities.Response {
  414. err := db.DeleteByCondition(
  415. models.InstallTask{
  416. Model: models.Model{
  417. ID: task_id,
  418. },
  419. TenantID: tenant_id,
  420. },
  421. )
  422. if err != nil {
  423. return exception.InternalServerError(err).ToResponse()
  424. }
  425. return entities.NewSuccessResponse(true)
  426. }
  427. func DeleteAllPluginInstallationTasks(
  428. tenant_id string,
  429. ) *entities.Response {
  430. err := db.DeleteByCondition(
  431. models.InstallTask{
  432. TenantID: tenant_id,
  433. },
  434. )
  435. if err != nil {
  436. return exception.InternalServerError(err).ToResponse()
  437. }
  438. return entities.NewSuccessResponse(true)
  439. }
  440. func DeletePluginInstallationItemFromTask(
  441. tenant_id string,
  442. task_id string,
  443. identifier plugin_entities.PluginUniqueIdentifier,
  444. ) *entities.Response {
  445. err := db.WithTransaction(func(tx *gorm.DB) error {
  446. item, err := db.GetOne[models.InstallTask](
  447. db.WithTransactionContext(tx),
  448. db.Equal("id", task_id),
  449. db.Equal("tenant_id", tenant_id),
  450. db.WLock(),
  451. )
  452. if err != nil {
  453. return err
  454. }
  455. plugins := []models.InstallTaskPluginStatus{}
  456. for _, plugin := range item.Plugins {
  457. if plugin.PluginUniqueIdentifier != identifier {
  458. plugins = append(plugins, plugin)
  459. }
  460. }
  461. successes := 0
  462. for _, plugin := range plugins {
  463. if plugin.Status == models.InstallTaskStatusSuccess {
  464. successes++
  465. }
  466. }
  467. if len(plugins) == successes {
  468. // delete the task if all plugins are installed successfully
  469. err = db.Delete(&item, tx)
  470. } else {
  471. item.Plugins = plugins
  472. err = db.Update(&item, tx)
  473. }
  474. return err
  475. })
  476. if err != nil {
  477. return exception.InternalServerError(err).ToResponse()
  478. }
  479. return entities.NewSuccessResponse(true)
  480. }
  481. func FetchPluginFromIdentifier(
  482. pluginUniqueIdentifier plugin_entities.PluginUniqueIdentifier,
  483. ) *entities.Response {
  484. _, err := db.GetOne[models.Plugin](
  485. db.Equal("plugin_unique_identifier", pluginUniqueIdentifier.String()),
  486. )
  487. if err == db.ErrDatabaseNotFound {
  488. return entities.NewSuccessResponse(false)
  489. }
  490. if err != nil {
  491. return exception.InternalServerError(err).ToResponse()
  492. }
  493. return entities.NewSuccessResponse(true)
  494. }
  495. func UninstallPlugin(
  496. tenant_id string,
  497. plugin_installation_id string,
  498. ) *entities.Response {
  499. // Check if the plugin exists for the tenant
  500. installation, err := db.GetOne[models.PluginInstallation](
  501. db.Equal("tenant_id", tenant_id),
  502. db.Equal("id", plugin_installation_id),
  503. )
  504. if err == db.ErrDatabaseNotFound {
  505. return exception.ErrPluginNotFound().ToResponse()
  506. }
  507. if err != nil {
  508. return exception.InternalServerError(err).ToResponse()
  509. }
  510. pluginUniqueIdentifier, err := plugin_entities.NewPluginUniqueIdentifier(installation.PluginUniqueIdentifier)
  511. if err != nil {
  512. return exception.UniqueIdentifierError(err).ToResponse()
  513. }
  514. // get declaration
  515. declaration, err := helper.CombinedGetPluginDeclaration(
  516. pluginUniqueIdentifier,
  517. plugin_entities.PluginRuntimeType(installation.RuntimeType),
  518. )
  519. if err != nil {
  520. return exception.InternalServerError(err).ToResponse()
  521. }
  522. // Uninstall the plugin
  523. deleteResponse, err := curd.UninstallPlugin(
  524. tenant_id,
  525. pluginUniqueIdentifier,
  526. installation.ID,
  527. declaration,
  528. )
  529. if err != nil {
  530. return exception.InternalServerError(fmt.Errorf("failed to uninstall plugin: %s", err.Error())).ToResponse()
  531. }
  532. if deleteResponse.IsPluginDeleted {
  533. // delete the plugin if no installation left
  534. manager := plugin_manager.Manager()
  535. if deleteResponse.Installation.RuntimeType == string(
  536. plugin_entities.PLUGIN_RUNTIME_TYPE_LOCAL,
  537. ) {
  538. err = manager.UninstallFromLocal(pluginUniqueIdentifier)
  539. if err != nil {
  540. return exception.InternalServerError(fmt.Errorf("failed to uninstall plugin: %s", err.Error())).ToResponse()
  541. }
  542. }
  543. }
  544. return entities.NewSuccessResponse(true)
  545. }