@cryptotaxi247 / netdata-1 / commits / 9be211256

chore: remove lock files from go.d/python.d (#19668)

* remove locks from go.d * remove lock files from pythond

Ilya Mashchenko committed Feb 18, 2025 at 11:49 UTC 9be211256c26ce13ea62aa94fe6f4e58bfe3b119
7 files changed +3 -131
src/collectors/python.d.plugin/python.d.plugin.in
+3 -79
@@ -51,7 +51,6 @@ ENV_NETDATA_PLUGINS_DIR = 'NETDATA_PLUGINS_DIR'
51 ENV_NETDATA_USER_PLUGINS_DIRS = 'NETDATA_USER_PLUGINS_DIRS'
52 ENV_NETDATA_LIB_DIR = 'NETDATA_LIB_DIR'
53 ENV_NETDATA_UPDATE_EVERY = 'NETDATA_UPDATE_EVERY'
54 -ENV_NETDATA_LOCK_DIR = 'NETDATA_LOCK_DIR'
54
55
56 def add_pythond_packages():
@@ -66,7 +65,6 @@ add_pythond_packages()
65 from bases.collection import safe_print
66 from bases.loggers import PythonDLogger
67 from bases.loaders import load_config
69 -from third_party import filelock
68
69 try:
70 from collections import OrderedDict
@@ -91,10 +89,7 @@ def dirs():
89 ENV_NETDATA_PLUGINS_DIR,
90 os.path.dirname(__file__),
91 )
94 - locks = os.getenv(
95 - ENV_NETDATA_LOCK_DIR,
96 - os.path.join('@varlibdir_POST@', 'lock')
97 - )
92 +
93 modules_user_config = os.path.join(plugin_user_config, 'python.d')
94 modules_stock_config = os.path.join(plugin_stock_config, 'python.d')
95 modules = os.path.abspath(pluginsd + '/../python.d')
@@ -112,7 +107,6 @@ def dirs():
107 'modules',
108 'user_modules',
109 'var_lib',
115 - 'locks',
110 ]
111 )
112 return Dirs(
@@ -123,7 +117,6 @@ def dirs():
117 modules,
118 user_modules,
119 var_lib,
126 - locks,
120 )
121
122
@@ -501,54 +494,15 @@ class PluginConfig(dict):
494 return self['default_run']
495
496
504 -class FileLockRegistry:
505 - def __init__(self, path):
506 - self.path = path
507 - self.locks = dict()
508 -
509 - @staticmethod
510 - def rename(name):
511 - # go version name is 'docker'
512 - if name.startswith("dockerd"):
513 - name = "docker" + name[7:]
514 - return name
515 -
516 - def register(self, name):
517 - name = self.rename(name)
518 - if name in self.locks:
519 - return
520 - file = os.path.join(self.path, '{0}.collector.lock'.format(name))
521 - lock = filelock.FileLock(file)
522 - lock.acquire(timeout=0)
523 - self.locks[name] = lock
524 -
525 - def unregister(self, name):
526 - name = self.rename(name)
527 - if name not in self.locks:
528 - return
529 - lock = self.locks[name]
530 - lock.release()
531 - del self.locks[name]
532 -
533 -
534 -class DummyRegistry:
535 - def register(self, name):
536 - pass
537 -
538 - def unregister(self, name):
539 - pass
540 -
541 -
497 class Plugin:
498 config_name = 'python.d.conf'
499 jobs_status_dump_name = 'pythond-jobs-statuses.json'
500
546 - def __init__(self, modules_to_run, min_update_every, registry):
501 + def __init__(self, modules_to_run, min_update_every):
502 self.modules_to_run = modules_to_run
503 self.min_update_every = min_update_every
504 self.config = PluginConfig(PLUGIN_BASE_CONF)
505 self.log = PythonDLogger()
551 - self.registry = registry
506 self.started_jobs = collections.defaultdict(dict)
507 self.jobs = list()
508 self.saver = None
@@ -706,30 +660,12 @@ class Plugin:
660 continue
661 self.log.info('{0}[{1}] : check success'.format(job.module_name, job.real_name))
662
709 - try:
710 - self.registry.register(job.full_name())
711 - except filelock.Timeout as error:
712 - self.log.info('{0}[{1}] : already registered by another process, skipping the job ({2})'.format(
713 - job.module_name, job.real_name, error))
714 - job.status = JOB_STATUS_DROPPED
715 - continue
716 - except Exception as error:
717 - self.log.warning('{0}[{1}] : registration failed: {2}, skipping the job'.format(
718 - job.module_name, job.real_name, error))
719 - job.status = JOB_STATUS_DROPPED
720 - continue
721 -
663 try:
664 job.create()
665 except Exception as error:
666 self.log.warning("{0}[{1}] : unhandled exception on create : {2}, skipping the job".format(
667 job.module_name, job.real_name, repr(error)))
668 job.status = JOB_STATUS_DROPPED
728 - try:
729 - self.registry.unregister(job.full_name())
730 - except Exception as error:
731 - self.log.warning('{0}[{1}] : deregistration failed: {2}'.format(
732 - job.module_name, job.real_name, error))
669 continue
670
671 self.started_jobs[job.module_name] = job.actual_name
@@ -803,7 +739,6 @@ def parse_command_line():
739
740 debug = False
741 trace = False
806 - nolock = False
742 update_every = 1
743 modules_to_run = list()
744
@@ -820,9 +755,6 @@ def parse_command_line():
755 if 'trace' in opts:
756 trace = True
757 opts.remove('trace')
823 - if 'nolock' in opts:
824 - nolock = True
825 - opts.remove('nolock')
758 if opts:
759 modules_to_run = list(opts)
760
@@ -832,14 +764,12 @@ def parse_command_line():
764 'update_every',
765 'debug',
766 'trace',
835 - 'nolock',
767 'modules_to_run',
768 ])
769 return cmd(
770 update_every,
771 debug,
772 trace,
842 - nolock,
773 modules_to_run,
774 )
775
@@ -907,11 +837,6 @@ def main():
837
838 log.info('using python v{0}'.format(PY_VERSION[0]))
839
910 - if DIRS.locks and not cmd.nolock:
911 - registry = FileLockRegistry(DIRS.locks)
912 - else:
913 - registry = DummyRegistry()
914 -
840 unique_avail_module_names = set([m.name for m in AVAILABLE_MODULES])
841 unknown = set(cmd.modules_to_run) - unique_avail_module_names
842 if unknown:
@@ -924,8 +849,7 @@ def main():
849 p = Plugin(
850 get_modules_to_run(cmd),
851 cmd.update_every,
927 - registry,
928 - )
852 +)
853
854 # cheap attempt to reduce chance of python.d job running before go.d
855 # TODO: better implementation needed
src/go/cmd/godplugin/config.go
-5
@@ -19,7 +19,6 @@ type envConfig struct {
19 userDir string
20 stockDir string
21 varLibDir string
22 - lockDir string
22 watchPath string
23 logLevel string
24 }
@@ -30,7 +29,6 @@ func newEnvConfig() *envConfig {
29 userDir: os.Getenv("NETDATA_USER_CONFIG_DIR"),
30 stockDir: os.Getenv("NETDATA_STOCK_CONFIG_DIR"),
31 varLibDir: os.Getenv("NETDATA_LIB_DIR"),
33 - lockDir: os.Getenv("NETDATA_LOCK_DIR"),
32 watchPath: os.Getenv("NETDATA_PLUGINS_GOD_WATCH_PATH"),
33 logLevel: os.Getenv("NETDATA_LOG_LEVEL"),
34 }
@@ -38,7 +36,6 @@ func newEnvConfig() *envConfig {
36 cfg.userDir = cfg.handleDirOnWin(cfg.userDir)
37 cfg.stockDir = cfg.handleDirOnWin(cfg.stockDir)
38 cfg.varLibDir = cfg.handleDirOnWin(cfg.varLibDir)
41 - cfg.lockDir = cfg.handleDirOnWin(cfg.lockDir)
39 cfg.watchPath = cfg.handleDirOnWin(cfg.watchPath)
40
41 return cfg
@@ -66,7 +63,6 @@ type config struct {
63 collectorsWatchPath []string
64 serviceDiscoveryDir multipath.MultiPath
65 stateFile string
69 - lockDir string
66 }
67
68 func newConfig(opts *cli.Option, env *envConfig) *config {
@@ -79,7 +75,6 @@ func newConfig(opts *cli.Option, env *envConfig) *config {
75 cfg.collectorsWatchPath = cfg.initCollectorsWatchPaths(opts, env)
76 cfg.serviceDiscoveryDir = cfg.initServiceDiscoveryConfigDir()
77 cfg.stateFile = cfg.initStateFile(env)
82 - cfg.lockDir = env.lockDir
78
79 return cfg
80 }
src/go/cmd/godplugin/main.go
-1
@@ -54,7 +54,6 @@ func main() {
54 ServiceDiscoveryConfigDir: cfg.serviceDiscoveryDir,
55 CollectorsConfigWatchPath: cfg.collectorsWatchPath,
56 StateFile: cfg.stateFile,
57 - LockDir: cfg.lockDir,
57 RunModule: opts.Module,
58 MinUpdateEvery: opts.UpdateEvery,
59 })
src/go/plugin/go.d/agent/agent.go
-8
@@ -18,7 +18,6 @@ import (
18 "github.com/netdata/netdata/go/plugins/pkg/safewriter"
19 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/confgroup"
20 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/discovery"
21 - "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/filelock"
21 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/filestatus"
22 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/functions"
23 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/jobmgr"
@@ -37,7 +36,6 @@ type Config struct {
36 CollectorsConfigWatchPath []string
37 ServiceDiscoveryConfigDir []string
38 StateFile string
40 - LockDir string
39 ModuleRegistry module.Registry
40 RunModule string
41 MinUpdateEvery int
@@ -55,7 +53,6 @@ type Agent struct {
53 ServiceDiscoveryConfigDir multipath.MultiPath
54
55 StateFile string
58 - LockDir string
56
57 RunModule string
58 MinUpdateEvery int
@@ -80,7 +77,6 @@ func New(cfg Config) *Agent {
77 ServiceDiscoveryConfigDir: cfg.ServiceDiscoveryConfigDir,
78 CollectorsConfigWatchPath: cfg.CollectorsConfigWatchPath,
79 StateFile: cfg.StateFile,
83 - LockDir: cfg.LockDir,
80 RunModule: cfg.RunModule,
81 MinUpdateEvery: cfg.MinUpdateEvery,
82 ModuleRegistry: module.DefaultRegistry,
@@ -205,10 +201,6 @@ func (a *Agent) run(ctx context.Context) {
201 jobMgr.Vnodes = reg
202 }
203
208 - if a.LockDir != "" {
209 - jobMgr.FileLock = filelock.New(a.LockDir)
210 - }
211 -
204 var fsMgr *filestatus.Manager
205 if !isTerminal && a.StateFile != "" {
206 fsMgr = filestatus.NewManager(a.StateFile)
src/go/plugin/go.d/agent/jobmgr/di.go
-6
@@ -9,12 +9,6 @@ import (
9 "github.com/netdata/netdata/go/plugins/plugin/go.d/agent/vnodes"
10 )
11
12 -type FileLocker interface {
13 - Lock(name string) (bool, error)
14 - Unlock(name string)
15 - UnlockAll()
16 -}
17 -
12 type FileStatus interface {
13 Save(cfg confgroup.Config, state string)
14 Remove(cfg confgroup.Config)
src/go/plugin/go.d/agent/jobmgr/dyncfg_collector.go
-27
@@ -342,7 +342,6 @@ func (m *Manager) dyncfgConfigRestart(fn functions.Function) {
342 return
343 case dyncfgRunning:
344 m.FileStatus.Remove(ecfg.cfg)
345 - m.FileLock.Unlock(ecfg.cfg.FullName())
345 m.stopRunningJob(ecfg.cfg.FullName())
346 default:
347 }
@@ -355,14 +354,6 @@ func (m *Manager) dyncfgConfigRestart(fn functions.Function) {
354 return
355 }
356
358 - if ok, err := m.FileLock.Lock(ecfg.cfg.FullName()); !ok && err == nil {
359 - job.Cleanup()
360 - ecfg.status = dyncfgFailed
361 - m.dyncfgRespf(fn, 500, "Job restart failed: cannot filelock.")
362 - m.dyncfgJobStatus(ecfg.cfg, ecfg.status)
363 - return
364 - }
365 -
357 ecfg.status = dyncfgRunning
358
359 if isDyncfg(ecfg.cfg) {
@@ -439,14 +430,6 @@ func (m *Manager) dyncfgConfigEnable(fn functions.Function) {
430 return
431 }
432
442 - if ok, err := m.FileLock.Lock(ecfg.cfg.FullName()); !ok && err == nil {
443 - job.Cleanup()
444 - ecfg.status = dyncfgFailed
445 - m.dyncfgRespf(fn, 500, "Job enable failed: can not filelock.")
446 - m.dyncfgJobStatus(ecfg.cfg, ecfg.status)
447 - return
448 - }
449 -
433 ecfg.status = dyncfgRunning
434
435 if isDyncfg(ecfg.cfg) {
@@ -489,7 +472,6 @@ func (m *Manager) dyncfgConfigDisable(fn functions.Function) {
472 if isDyncfg(ecfg.cfg) {
473 m.FileStatus.Remove(ecfg.cfg)
474 }
492 - m.FileLock.Unlock(ecfg.cfg.FullName())
475 default:
476 }
477
@@ -583,7 +565,6 @@ func (m *Manager) dyncfgConfigRemove(fn functions.Function) {
565 m.seenConfigs.remove(ecfg.cfg)
566 m.exposedConfigs.remove(ecfg.cfg)
567 m.stopRunningJob(ecfg.cfg.FullName())
586 - m.FileLock.Unlock(ecfg.cfg.FullName())
568 m.FileStatus.Remove(ecfg.cfg)
569
570 m.dyncfgRespf(fn, 200, "")
@@ -666,14 +647,6 @@ func (m *Manager) dyncfgConfigUpdate(fn functions.Function) {
647 return
648 }
649
669 - if ok, err := m.FileLock.Lock(scfg.cfg.FullName()); !ok && err == nil {
670 - job.Cleanup()
671 - scfg.status = dyncfgFailed
672 - m.dyncfgRespf(fn, 500, "Job update failed: cannot create file lock.")
673 - m.dyncfgJobStatus(scfg.cfg, scfg.status)
674 - return
675 - }
676 -
650 scfg.status = dyncfgRunning
651 m.startRunningJob(job)
652 m.dyncfgRespf(fn, 200, "")
src/go/plugin/go.d/agent/jobmgr/manager.go
-5
@@ -33,7 +33,6 @@ func New() *Manager {
33 slog.String("component", "job manager"),
34 ),
35 Out: io.Discard,
36 - FileLock: noop{},
36 FileStatus: noop{},
37 FileStatusStore: noop{},
38 FnReg: noop{},
@@ -64,7 +63,6 @@ type Manager struct {
63 Modules module.Registry
64 ConfigDefaults confgroup.Registry
65
67 - FileLock FileLocker
66 FileStatus FileStatus
67 FileStatusStore FileStatusStore
68 FnReg FunctionRegistry
@@ -193,7 +191,6 @@ func (m *Manager) addConfig(cfg confgroup.Config) {
191 }
192 if ecfg.status == dyncfgRunning {
193 m.stopRunningJob(ecfg.cfg.FullName())
196 - m.FileLock.Unlock(ecfg.cfg.FullName())
194 m.FileStatus.Remove(ecfg.cfg)
195 }
196 scfg.status = dyncfgAccepted
@@ -226,7 +223,6 @@ func (m *Manager) removeConfig(cfg confgroup.Config) {
223
224 m.exposedConfigs.remove(cfg)
225 m.stopRunningJob(cfg.FullName())
229 - m.FileLock.Unlock(cfg.FullName())
226 m.FileStatus.Remove(cfg)
227
228 if !isStock(cfg) || ecfg.status == dyncfgRunning {
@@ -273,7 +269,6 @@ func (m *Manager) stopRunningJob(name string) {
269 }
270
271 func (m *Manager) cleanup() {
276 - m.FileLock.UnlockAll()
272 m.FnReg.Unregister("config")
273
274 m.runningJobs.lock()