diff --git a/plugin/environment.go b/plugin/environment.go index 5114474fc8..8a63c89cdb 100644 --- a/plugin/environment.go +++ b/plugin/environment.go @@ -107,10 +107,11 @@ func (env *Environment) PrepackagedPlugins() []*PrepackagedPlugin { } // Returns a list of all currently active plugins within the environment. +// The returned list should not be modified. func (env *Environment) Active() []*model.BundleInfo { activePlugins := []*model.BundleInfo{} env.registeredPlugins.Range(func(key, value interface{}) bool { - plugin := value.(*registeredPlugin) + plugin := value.(registeredPlugin) if env.IsActive(plugin.BundleInfo.Manifest.Id) { activePlugins = append(activePlugins, plugin.BundleInfo) } @@ -133,13 +134,15 @@ func (env *Environment) GetPluginState(id string) int { return model.PluginStateNotRunning } - return rp.(*registeredPlugin).State + return rp.(registeredPlugin).State } // SetPluginState sets the current state of a plugin (disabled, running, or error) func (env *Environment) SetPluginState(id string, state int) { if rp, ok := env.registeredPlugins.Load(id); ok { - rp.(*registeredPlugin).State = state + p := rp.(registeredPlugin) + p.State = state + env.registeredPlugins.Store(id, p) } } @@ -226,13 +229,12 @@ func (env *Environment) Activate(id string) (manifest *model.Manifest, activated value, ok := env.registeredPlugins.Load(id) if !ok { value = newRegisteredPlugin(pluginInfo) - env.registeredPlugins.Store(id, value) } - rp := value.(*registeredPlugin) - + rp := value.(registeredPlugin) // Store latest BundleInfo in case something has changed since last activation rp.BundleInfo = pluginInfo + env.registeredPlugins.Store(id, rp) defer func() { if reterr == nil { @@ -270,6 +272,7 @@ func (env *Environment) Activate(id string) (manifest *model.Manifest, activated return nil, false, errors.Wrapf(err, "unable to start plugin: %v", id) } rp.supervisor = sup + env.registeredPlugins.Store(id, rp) componentActivated = true } @@ -302,7 +305,7 @@ func (env *Environment) Deactivate(id string) bool { return false } - rp := p.(*registeredPlugin) + rp := p.(registeredPlugin) if rp.supervisor != nil { if err := rp.supervisor.Hooks().OnDeactivate(); err != nil { env.logger.Error("Plugin OnDeactivate() error", mlog.String("plugin_id", rp.BundleInfo.Manifest.Id), mlog.Err(err)) @@ -328,7 +331,7 @@ func (env *Environment) Shutdown() { var wg sync.WaitGroup env.registeredPlugins.Range(func(key, value interface{}) bool { - rp := value.(*registeredPlugin) + rp := value.(registeredPlugin) if rp.supervisor == nil { return true @@ -430,7 +433,7 @@ func (env *Environment) UnpackWebappBundle(id string) (*model.Manifest, error) { // Consider using RunMultiPluginHook instead. func (env *Environment) HooksForPlugin(id string) (Hooks, error) { if p, ok := env.registeredPlugins.Load(id); ok { - rp := p.(*registeredPlugin) + rp := p.(registeredPlugin) if rp.supervisor != nil { return rp.supervisor.Hooks(), nil } @@ -445,7 +448,7 @@ func (env *Environment) HooksForPlugin(id string) (Hooks, error) { // plugins is not specified. func (env *Environment) RunMultiPluginHook(hookRunnerFunc func(hooks Hooks) bool, hookId int) { env.registeredPlugins.Range(func(key, value interface{}) bool { - rp := value.(*registeredPlugin) + rp := value.(registeredPlugin) if rp.supervisor == nil || !rp.supervisor.Implements(hookId) { return true @@ -465,7 +468,7 @@ func (env *Environment) SetPrepackagedPlugins(plugins []*PrepackagedPlugin) { env.prepackagedPluginsLock.Unlock() } -func newRegisteredPlugin(bundle *model.BundleInfo) *registeredPlugin { +func newRegisteredPlugin(bundle *model.BundleInfo) registeredPlugin { state := model.PluginStateNotRunning - return ®isteredPlugin{failTimeStamps: []time.Time{}, State: state, BundleInfo: bundle} + return registeredPlugin{failTimeStamps: []time.Time{}, State: state, BundleInfo: bundle} } diff --git a/plugin/health_check.go b/plugin/health_check.go index 0c873f82e4..24e533f672 100644 --- a/plugin/health_check.go +++ b/plugin/health_check.go @@ -77,7 +77,7 @@ func (job *PluginHealthCheckJob) checkPlugin(id string) { if !ok { return } - rp := p.(*registeredPlugin) + rp := p.(registeredPlugin) sup := rp.supervisor if sup == nil { @@ -98,14 +98,16 @@ func (job *PluginHealthCheckJob) handleHealthCheckFail(id string, err error) { if !ok { return } - p := rp.(*registeredPlugin) + p := rp.(registeredPlugin) // Append current failure before checking for deactivate vs restart action p.failTimeStamps = append(p.failTimeStamps, time.Now()) p.lastError = err + job.env.registeredPlugins.Store(id, p) if shouldDeactivatePlugin(p) { p.failTimeStamps = []time.Time{} + job.env.registeredPlugins.Store(id, p) mlog.Debug("Deactivating plugin due to multiple crashes", mlog.String("id", id)) job.env.Deactivate(id) job.env.SetPluginState(id, model.PluginStateFailedToStayRunning) @@ -135,7 +137,7 @@ func (job *PluginHealthCheckJob) Cancel() { // shouldDeactivatePlugin determines if a plugin needs to be deactivated after certain criteria is met. // // The criteria is based on if the plugin has consistently failed during the configured number of restarts, within the configured time window. -func shouldDeactivatePlugin(rp *registeredPlugin) bool { +func shouldDeactivatePlugin(rp registeredPlugin) bool { if len(rp.failTimeStamps) >= HEALTH_CHECK_RESTART_LIMIT { index := len(rp.failTimeStamps) - HEALTH_CHECK_RESTART_LIMIT t := rp.failTimeStamps[index] diff --git a/plugin/supervisor.go b/plugin/supervisor.go index 3b5ecba8d6..824e5966d3 100644 --- a/plugin/supervisor.go +++ b/plugin/supervisor.go @@ -9,6 +9,7 @@ import ( "path/filepath" "runtime" "strings" + "sync" "time" plugin "github.com/hashicorp/go-plugin" @@ -17,6 +18,7 @@ import ( ) type supervisor struct { + lock sync.RWMutex client *plugin.Client hooks Hooks implemented [TotalHooksId]bool @@ -99,17 +101,22 @@ func newSupervisor(pluginInfo *model.BundleInfo, parentLogger *mlog.Logger, apiI } func (sup *supervisor) Shutdown() { + sup.lock.RLock() + defer sup.lock.RUnlock() if sup.client != nil { sup.client.Kill() } } func (sup *supervisor) Hooks() Hooks { + sup.lock.RLock() + defer sup.lock.RUnlock() return sup.hooks } // PerformHealthCheck checks the plugin through an an RPC ping. func (sup *supervisor) PerformHealthCheck() error { + // No need for a lock here because Ping is read-locked. if pingErr := sup.Ping(); pingErr != nil { for pingFails := 1; pingFails < HEALTH_CHECK_PING_FAIL_LIMIT; pingFails++ { pingErr = sup.Ping() @@ -128,8 +135,9 @@ func (sup *supervisor) PerformHealthCheck() error { // Ping checks that the RPC connection with the plugin is alive and healthy. func (sup *supervisor) Ping() error { + sup.lock.RLock() + defer sup.lock.RUnlock() client, err := sup.client.Client() - if err != nil { return err } @@ -138,5 +146,7 @@ func (sup *supervisor) Ping() error { } func (sup *supervisor) Implements(hookId int) bool { + sup.lock.RLock() + defer sup.lock.RUnlock() return sup.implemented[hookId] }