From 693d3482835134eb2fca355789c32e6e0ff2f3c0 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Fri, 29 Nov 2019 17:47:29 -0800 Subject: [PATCH 1/7] what is needed here is only string value Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index 58b2923ef9..a0b0c29f55 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -210,11 +210,11 @@ def kill_worker end def supervisor_dump_config_handler - $log.info config[:fluentd_conf].to_s + $log.info config[:fluentd_conf] end def supervisor_get_dump_config_handler - {conf: config[:fluentd_conf].to_s} + {conf: config[:fluentd_conf]} end end @@ -322,7 +322,7 @@ def self.load_config(path, params = {}) path, JSON.dump(params)], command_sender: command_sender, - fluentd_conf: fluentd_conf, + fluentd_conf: fluentd_conf.to_s, main_cmd: main_cmd, signame: signame, } From 5e2f2c9f2c56ca7bdb7a849faf3bdab5e499d501 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Wed, 27 Nov 2019 16:04:38 -0800 Subject: [PATCH 2/7] Pass original config file Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index a0b0c29f55..0247876080 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -322,7 +322,7 @@ def self.load_config(path, params = {}) path, JSON.dump(params)], command_sender: command_sender, - fluentd_conf: fluentd_conf.to_s, + fluentd_conf: params['fluentd_conf'], main_cmd: main_cmd, signame: signame, } @@ -640,6 +640,7 @@ def supervise 'use_v1_config' => @use_v1_config, 'conf_encoding' => @conf_encoding, 'signame' => @signame, + 'fluentd_conf' => @conf.to_s, 'workers' => @system_config.workers, 'root_dir' => @system_config.root_dir, From 2f956e41a08b161372eaed7522b9dc33baf3d388 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Wed, 27 Nov 2019 14:55:24 -0800 Subject: [PATCH 3/7] Pass system_config parameters at #supervise Current implementation loads system_config parameters(rpc_endpoint, enable_get_dump, counter_server) in Supervisor.load_config which is called at ServerEngine creation time. And it can also reload config dynamically via serverengine's mechanism(sending signal USR2 to fluentd). These parameters can be reloadable but it doesn't make any changes to fluentd itself. fluentd document also doesn't mention about signal USR2. so it's no need to load in Supervisor.load_config. it's enough to pass the value at start time. Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 21 ++++++--------------- 1 file changed, 6 insertions(+), 15 deletions(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index 0247876080..7021a2087a 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -234,7 +234,6 @@ def after_start class Supervisor def self.load_config(path, params = {}) - pre_loadtime = 0 pre_loadtime = params['pre_loadtime'].to_i if params['pre_loadtime'] pre_config_mtime = nil @@ -246,17 +245,6 @@ def self.load_config(path, params = {}) return params['pre_conf'] end - config_fname = File.basename(path) - config_basedir = File.dirname(path) - # Assume fluent.conf encoding is UTF-8 - config_data = File.open(path, "r:#{params['conf_encoding']}:utf-8") {|f| f.read } - inline_config = params['inline_config'] - if inline_config - config_data << "\n" << inline_config.gsub("\\n","\n") - end - fluentd_conf = Fluent::Config.parse(config_data, config_fname, config_basedir, params['use_v1_config']) - system_config = SystemConfig.create(fluentd_conf) - # these params must NOT be configured via system config here. # these may be overridden by command line params. workers = params['workers'] @@ -269,9 +257,9 @@ def self.load_config(path, params = {}) chgroup = params['chgroup'] log_rotate_age = params['log_rotate_age'] log_rotate_size = params['log_rotate_size'] - rpc_endpoint = system_config.rpc_endpoint - enable_get_dump = system_config.enable_get_dump - counter_server = system_config.counter_server + rpc_endpoint = params['rpc_endpoint'] + enable_get_dump = params['enable_get_dump'] + counter_server = params['counter_server'] log_opts = {suppress_repeated_stacktrace: suppress_repeated_stacktrace} logger_initializer = Supervisor::LoggerInitializer.new( @@ -646,6 +634,9 @@ def supervise 'root_dir' => @system_config.root_dir, 'log_level' => @system_config.log_level, 'suppress_repeated_stacktrace' => @system_config.suppress_repeated_stacktrace, + 'rpc_endpoint' => @system_config.rpc_endpoint, + 'enable_get_dump' => @system_config.enable_get_dump, + 'counter_server' => @system_config.counter_server, } se = ServerEngine.create(ServerModule, WorkerModule){ From 1ac1149315db0b2d0cb88d871e0da78745e7d784 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Wed, 27 Nov 2019 15:01:00 -0800 Subject: [PATCH 4/7] delete old comment Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 2 -- 1 file changed, 2 deletions(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index 7021a2087a..304945cce0 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -245,8 +245,6 @@ def self.load_config(path, params = {}) return params['pre_conf'] end - # these params must NOT be configured via system config here. - # these may be overridden by command line params. workers = params['workers'] root_dir = params['root_dir'] log_level = params['log_level'] From 41178007939bafed47384baf193245c8e1fbb608 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Wed, 27 Nov 2019 15:03:29 -0800 Subject: [PATCH 5/7] Remove useless local variables Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 23 ++++++++--------------- 1 file changed, 8 insertions(+), 15 deletions(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index 304945cce0..11ab2bee0b 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -245,8 +245,6 @@ def self.load_config(path, params = {}) return params['pre_conf'] end - workers = params['workers'] - root_dir = params['root_dir'] log_level = params['log_level'] suppress_repeated_stacktrace = params['suppress_repeated_stacktrace'] @@ -255,9 +253,6 @@ def self.load_config(path, params = {}) chgroup = params['chgroup'] log_rotate_age = params['log_rotate_age'] log_rotate_size = params['log_rotate_size'] - rpc_endpoint = params['rpc_endpoint'] - enable_get_dump = params['enable_get_dump'] - counter_server = params['counter_server'] log_opts = {suppress_repeated_stacktrace: suppress_repeated_stacktrace} logger_initializer = Supervisor::LoggerInitializer.new( @@ -274,12 +269,10 @@ def self.load_config(path, params = {}) # ServerEngine's "daemonize" option is boolean, and path of pid file is brought by "pid_path" pid_path = params['daemonize'] daemonize = !!params['daemonize'] - main_cmd = params['main_cmd'] - signame = params['signame'] se_config = { worker_type: 'spawn', - workers: workers, + workers: params['workers'], log_stdin: false, log_stdout: false, log_stderr: false, @@ -287,7 +280,7 @@ def self.load_config(path, params = {}) auto_heartbeat: false, unrecoverable_exit_codes: [2], stop_immediately_at_unrecoverable_exit: true, - root_dir: root_dir, + root_dir: params['root_dir'], logger: logger, log: logger.out, log_path: log_path, @@ -298,9 +291,9 @@ def self.load_config(path, params = {}) chumask: 0, suppress_repeated_stacktrace: suppress_repeated_stacktrace, daemonize: daemonize, - rpc_endpoint: rpc_endpoint, - counter_server: counter_server, - enable_get_dump: enable_get_dump, + rpc_endpoint: params['rpc_endpoint'], + counter_server: params['counter_server'], + enable_get_dump: params['enable_get_dump'], windows_daemon_cmdline: [ServerEngine.ruby_bin_path, File.join(File.dirname(__FILE__), 'daemon.rb'), ServerModule.name, @@ -309,8 +302,8 @@ def self.load_config(path, params = {}) JSON.dump(params)], command_sender: command_sender, fluentd_conf: params['fluentd_conf'], - main_cmd: main_cmd, - signame: signame, + main_cmd: params['main_cmd'], + signame: params['signame'], } if daemonize se_config[:pid_path] = pid_path @@ -323,7 +316,7 @@ def self.load_config(path, params = {}) pre_params['pre_conf'] = nil params['pre_conf'][:windows_daemon_cmdline][5] = JSON.dump(pre_params) - return se_config + se_config end class LoggerInitializer From b6a624423b227dfa1d8acee6924cb325356d28f8 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Wed, 27 Nov 2019 15:04:06 -0800 Subject: [PATCH 6/7] Fix indent Signed-off-by: Yuta Iwama --- lib/fluent/supervisor.rb | 66 ++++++++++++++++++++-------------------- 1 file changed, 33 insertions(+), 33 deletions(-) diff --git a/lib/fluent/supervisor.rb b/lib/fluent/supervisor.rb index 11ab2bee0b..12def0e7da 100644 --- a/lib/fluent/supervisor.rb +++ b/lib/fluent/supervisor.rb @@ -271,39 +271,39 @@ def self.load_config(path, params = {}) daemonize = !!params['daemonize'] se_config = { - worker_type: 'spawn', - workers: params['workers'], - log_stdin: false, - log_stdout: false, - log_stderr: false, - enable_heartbeat: true, - auto_heartbeat: false, - unrecoverable_exit_codes: [2], - stop_immediately_at_unrecoverable_exit: true, - root_dir: params['root_dir'], - logger: logger, - log: logger.out, - log_path: log_path, - log_level: log_level, - logger_initializer: logger_initializer, - chuser: chuser, - chgroup: chgroup, - chumask: 0, - suppress_repeated_stacktrace: suppress_repeated_stacktrace, - daemonize: daemonize, - rpc_endpoint: params['rpc_endpoint'], - counter_server: params['counter_server'], - enable_get_dump: params['enable_get_dump'], - windows_daemon_cmdline: [ServerEngine.ruby_bin_path, - File.join(File.dirname(__FILE__), 'daemon.rb'), - ServerModule.name, - WorkerModule.name, - path, - JSON.dump(params)], - command_sender: command_sender, - fluentd_conf: params['fluentd_conf'], - main_cmd: params['main_cmd'], - signame: params['signame'], + worker_type: 'spawn', + workers: params['workers'], + log_stdin: false, + log_stdout: false, + log_stderr: false, + enable_heartbeat: true, + auto_heartbeat: false, + unrecoverable_exit_codes: [2], + stop_immediately_at_unrecoverable_exit: true, + root_dir: params['root_dir'], + logger: logger, + log: logger.out, + log_path: log_path, + log_level: log_level, + logger_initializer: logger_initializer, + chuser: chuser, + chgroup: chgroup, + chumask: 0, + suppress_repeated_stacktrace: suppress_repeated_stacktrace, + daemonize: daemonize, + rpc_endpoint: params['rpc_endpoint'], + counter_server: params['counter_server'], + enable_get_dump: params['enable_get_dump'], + windows_daemon_cmdline: [ServerEngine.ruby_bin_path, + File.join(File.dirname(__FILE__), 'daemon.rb'), + ServerModule.name, + WorkerModule.name, + path, + JSON.dump(params)], + command_sender: command_sender, + fluentd_conf: params['fluentd_conf'], + main_cmd: params['main_cmd'], + signame: params['signame'], } if daemonize se_config[:pid_path] = pid_path From c98b348077536aeb8057d03c43ab3d00023d9787 Mon Sep 17 00:00:00 2001 From: Yuta Iwama Date: Fri, 29 Nov 2019 18:05:59 -0800 Subject: [PATCH 7/7] it's same as the test_read_config_with_multibyte_string Signed-off-by: Yuta Iwama --- test/test_supervisor.rb | 40 ---------------------------------------- 1 file changed, 40 deletions(-) diff --git a/test/test_supervisor.rb b/test/test_supervisor.rb index 3a50f9ac43..a8754b2f09 100644 --- a/test/test_supervisor.rb +++ b/test/test_supervisor.rb @@ -372,46 +372,6 @@ def test_load_config_for_daemonize assert_equal Fluent::Log::LEVEL_INFO, se_config[:log_level] end - def test_load_config_with_multibyte_string - tmp_path = "#{TMP_DIR}/dir/test_multibyte_config.conf" - conf_str = %[ - - @type forward - @id forward_input - @label @INPUT - - -] - FileUtils.mkdir_p(File.dirname(tmp_path)) - File.open(tmp_path, "w:utf-8") {|file| file.write(conf_str) } - - params = {} - params['workers'] = 1 - params['use_v1_config'] = true - params['log_path'] = 'test/tmp/supervisor/log' - params['suppress_repeated_stacktrace'] = true - params['log_level'] = Fluent::Log::LEVEL_INFO - params['conf_encoding'] = 'utf-8' - load_config_proc = Proc.new { Fluent::Supervisor.load_config(tmp_path, params) } - - se_config = load_config_proc.call - conf = se_config[:fluentd_conf] - label = conf.elements.detect {|e| e.name == "label" } - filter = label.elements.detect {|e| e.name == "filter" } - record_transformer = filter.elements.detect {|e| e.name = "record_transformer" } - assert_equal(Encoding::UTF_8, record_transformer["message"].encoding) - end - def test_logger opts = Fluent::Supervisor.default_options sv = Fluent::Supervisor.new(opts)