From cce83d9f97354eef13c4bcff500d6307c7d8980d Mon Sep 17 00:00:00 2001 From: HuaiYJ Date: Mon, 24 Aug 2026 16:21:14 +0800 Subject: [PATCH] fix: release plugin manager lock during scheduling --- internal/plugin/manager.go | 79 ++++++++++++++++++++++----------- internal/plugin/manager_test.go | 57 ++++++++++++++++++++++++ 2 files changed, 111 insertions(+), 25 deletions(-) diff --git a/internal/plugin/manager.go b/internal/plugin/manager.go index fecc8c3..096b5c5 100644 --- a/internal/plugin/manager.go +++ b/internal/plugin/manager.go @@ -37,6 +37,7 @@ type Manager struct { Cron *cron.Cron Lock sync.Mutex RunningPlugins map[string]*RunningPlugin + scheduling map[string]struct{} innerCancel context.CancelFunc log *zap.Logger @@ -51,6 +52,7 @@ func NewPluginManager(config *config.Config, log *zap.Logger) (*Manager, error) Plugins: make(map[string]*Plugin), Cron: cron.New(cron.WithParser(cron.NewParser(cron.Second | cron.Minute | cron.Hour | cron.Dom | cron.Month | cron.DowOptional))), RunningPlugins: make(map[string]*RunningPlugin), + scheduling: make(map[string]struct{}), log: log, } files, err := os.ReadDir(config.PluginsConfDir) @@ -112,31 +114,9 @@ func (m *Manager) Start() { // 添加所有插件的cron任务 for _, plugin := range m.Plugins { _, err := m.Cron.AddFunc(plugin.Cron, func() { - // 跑起来之前先检查一次是否有插件运行 - m.Lock.Lock() - defer m.Lock.Unlock() - rp, ok := m.RunningPlugins[plugin.Name] - if ok && rp.Plugin.Timeout < 0 { // 如果有正在执行的,且他的timeout小于0(常驻进程),则不进行启动 - m.log.Debug("Plugin is already running", zap.String("plugin", plugin.Name)) - return - } else if ok { - err := rp.Stop() - if err != nil { // 停止失败,这次不再操作,留待下次继续尝试 - m.log.Error("Failed to stop plugin", zap.String("plugin", plugin.Name), zap.Error(err)) - return - } - delete(m.RunningPlugins, plugin.Name) - } - - m.log.Info("Running plugin", zap.String("plugin", plugin.Name)) - rp, err := m.RunPlugin(plugin) // 拉起插件进程,如果进程内部处理完毕,则自动回收 - if err != nil { - m.log.Error("Failed to run plugin", zap.String("plugin", plugin.Name), zap.Error(err)) - return - } - - // 添加拉起的进程到runningPlugins中 - m.RunningPlugins[plugin.Name] = rp + m.runScheduledPlugin(plugin, + func(rp *RunningPlugin) error { return rp.Stop() }, + m.RunPlugin) }) if err != nil { m.log.Error("Failed to add plugin to cron", zap.String("plugin", plugin.Name), zap.Error(err)) @@ -151,6 +131,55 @@ func (m *Manager) Start() { m.LaunchPluginMonitor(ctxt) } +func (m *Manager) runScheduledPlugin( + plugin *Plugin, + stop func(*RunningPlugin) error, + run func(*Plugin) (*RunningPlugin, error), +) { + m.Lock.Lock() + if m.scheduling == nil { + m.scheduling = make(map[string]struct{}) + } + if _, busy := m.scheduling[plugin.Name]; busy { + m.Lock.Unlock() + return + } + existing, ok := m.RunningPlugins[plugin.Name] + if ok && existing.Plugin.Timeout < 0 { + m.Lock.Unlock() + m.log.Debug("Plugin is already running", zap.String("plugin", plugin.Name)) + return + } + m.scheduling[plugin.Name] = struct{}{} + if ok { + delete(m.RunningPlugins, plugin.Name) + } + m.Lock.Unlock() + + if ok { + if err := stop(existing); err != nil { + m.Lock.Lock() + m.RunningPlugins[plugin.Name] = existing + delete(m.scheduling, plugin.Name) + m.Lock.Unlock() + m.log.Error("Failed to stop plugin", zap.String("plugin", plugin.Name), zap.Error(err)) + return + } + } + + m.log.Info("Running plugin", zap.String("plugin", plugin.Name)) + running, err := run(plugin) + m.Lock.Lock() + delete(m.scheduling, plugin.Name) + if err == nil { + m.RunningPlugins[plugin.Name] = running + } + m.Lock.Unlock() + if err != nil { + m.log.Error("Failed to run plugin", zap.String("plugin", plugin.Name), zap.Error(err)) + } +} + // LaunchPluginMonitor 在后台 goroutine 中以 1 秒为周期调用 CleanupRunningPlugins, // ctxt 被取消后 goroutine 将退出。 func (m *Manager) LaunchPluginMonitor(ctxt context.Context) { diff --git a/internal/plugin/manager_test.go b/internal/plugin/manager_test.go index 53fa010..a371bca 100644 --- a/internal/plugin/manager_test.go +++ b/internal/plugin/manager_test.go @@ -49,6 +49,63 @@ run_on_boot: false } } +func TestRunScheduledPluginDoesNotHoldManagerLockDuringProcessOperations(t *testing.T) { + plugin := &Plugin{Name: "demo", Timeout: 60} + oldRunning := &RunningPlugin{Plugin: plugin} + newRunning := &RunningPlugin{Plugin: plugin} + mgr := &Manager{ + RunningPlugins: map[string]*RunningPlugin{"demo": oldRunning}, + scheduling: make(map[string]struct{}), + log: newTestLogger(), + } + stopEntered := make(chan struct{}) + releaseStop := make(chan struct{}) + runEntered := make(chan struct{}) + releaseRun := make(chan struct{}) + done := make(chan struct{}) + go func() { + mgr.runScheduledPlugin(plugin, + func(*RunningPlugin) error { + close(stopEntered) + <-releaseStop + return nil + }, + func(*Plugin) (*RunningPlugin, error) { + close(runEntered) + <-releaseRun + return newRunning, nil + }) + close(done) + }() + + assertLockAvailable := func(stage string) { + t.Helper() + locked := make(chan struct{}) + go func() { + mgr.Lock.Lock() + mgr.Lock.Unlock() + close(locked) + }() + select { + case <-locked: + case <-time.After(time.Second): + t.Fatalf("Manager.Lock held during %s", stage) + } + } + + <-stopEntered + assertLockAvailable("Stop") + close(releaseStop) + <-runEntered + assertLockAvailable("RunPlugin") + close(releaseRun) + <-done + + if got := mgr.RunningPlugins["demo"]; got != newRunning { + t.Fatalf("running plugin = %p, want %p", got, newRunning) + } +} + func TestNewPluginManager_SkipDisabledPlugin(t *testing.T) { confDir := t.TempDir() execDir := t.TempDir() -- Gitee