diff --git a/elasticsearch.go b/elasticsearch.go index e07bf9d..d9b8686 100644 --- a/elasticsearch.go +++ b/elasticsearch.go @@ -99,7 +99,7 @@ func (p *ESPlugin) Init(URI string) { p.done = make(chan bool) p.indexor.Start() - if Settings.Verbose { + if Settings.verbose() { // Only start the ErrorHandler goroutine when in verbose mode // no need to burn ressources otherwise go p.ErrorHandler() diff --git a/emitter.go b/emitter.go index 02d3e31..86a1fba 100644 --- a/emitter.go +++ b/emitter.go @@ -140,7 +140,7 @@ func CopyMulty(stop chan int, src io.Reader, writers ...io.Writer) error { payload := buf[:nr] meta := payloadMeta(payload) if len(meta) < 3 { - if Settings.Debug { + if Settings.debug() { Debug("[EMITTER] Found malformed record", string(payload[0:_maxN]), nr, "from:", src) } continue @@ -151,7 +151,7 @@ func CopyMulty(stop chan int, src io.Reader, writers ...io.Writer) error { log.Println("INFO: Large packet... We received ", len(payload), " bytes from ", src) } - if Settings.Debug { + if Settings.debug() { Debug("[EMITTER] input:", string(payload[0:_maxN]), nr, "from:", src) } @@ -172,7 +172,7 @@ func CopyMulty(stop chan int, src io.Reader, writers ...io.Writer) error { payload = append(payload[:headSize], body...) } - if Settings.Debug { + if Settings.debug() { Debug("[EMITTER] Rewritten input:", len(payload), "First 500 bytes:", string(payload[0:_maxN])) } } else { @@ -183,15 +183,15 @@ func CopyMulty(stop chan int, src io.Reader, writers ...io.Writer) error { } } - if Settings.PrettifyHTTP { + if Settings.prettifyHTTP() { payload = prettifyHTTP(payload) if len(payload) == 0 { continue } } - if Settings.SplitOutput { - if Settings.RecognizeTCPSessions { + if Settings.splitOutput() { + if Settings.recognizeTCPSessions() { if !PRO { log.Fatal("Detailed TCP sessions work only with PRO license") } diff --git a/emitter_test.go b/emitter_test.go index 9e2b9a5..2cbb9c7 100644 --- a/emitter_test.go +++ b/emitter_test.go @@ -32,7 +32,7 @@ func TestEmitter(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 1000; i++ { wg.Add(1) @@ -61,7 +61,7 @@ func TestEmitterFiltered(t *testing.T) { plugins.All = append(plugins.All, input, output) methods := HTTPMethods{[]byte("GET")} - Settings.ModifierConfig = HTTPModifierConfig{Methods: methods} + Settings.setModifierConfig(HTTPModifierConfig{Methods: methods}) emitter := &emitter{quit: quit} go emitter.Start(plugins, "") @@ -91,7 +91,7 @@ func TestEmitterFiltered(t *testing.T) { wg.Wait() emitter.Close() - Settings.ModifierConfig = HTTPModifierConfig{} + Settings.setModifierConfig(HTTPModifierConfig{}) } func TestEmitterSplitRoundRobin(t *testing.T) { @@ -117,10 +117,10 @@ func TestEmitterSplitRoundRobin(t *testing.T) { Outputs: []io.Writer{output1, output2}, } - Settings.SplitOutput = true + Settings.setSplitOutput(true) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 1000; i++ { wg.Add(1) @@ -135,7 +135,7 @@ func TestEmitterSplitRoundRobin(t *testing.T) { t.Errorf("Round robin should split traffic equally: %d vs %d", counter1, counter2) } - Settings.SplitOutput = false + Settings.setSplitOutput(false) } func TestEmitterRoundRobin(t *testing.T) { @@ -162,10 +162,10 @@ func TestEmitterRoundRobin(t *testing.T) { } plugins.All = append(plugins.All, input, output1, output2) - Settings.SplitOutput = true + Settings.setSplitOutput(true) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 1000; i++ { wg.Add(1) @@ -179,7 +179,7 @@ func TestEmitterRoundRobin(t *testing.T) { t.Errorf("Round robin should split traffic equally: %d vs %d", counter1, counter2) } - Settings.SplitOutput = false + Settings.setSplitOutput(false) } func TestEmitterSplitSession(t *testing.T) { @@ -222,11 +222,11 @@ func TestEmitterSplitSession(t *testing.T) { Outputs: []io.Writer{output1, output2}, } - Settings.SplitOutput = true - Settings.RecognizeTCPSessions = true + Settings.setSplitOutput(true) + Settings.setRecognizeTCPSessions(true) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 1000; i++ { // Keep session but randomize ACK @@ -247,8 +247,8 @@ func TestEmitterSplitSession(t *testing.T) { t.Errorf("Round robin should split traffic equally: %d vs %d", counter1, counter2) } - Settings.SplitOutput = false - Settings.RecognizeTCPSessions = false + Settings.setSplitOutput(false) + Settings.setRecognizeTCPSessions(false) emitter.Close() } @@ -269,7 +269,7 @@ func BenchmarkEmitter(b *testing.B) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) b.ResetTimer() diff --git a/gor.go b/gor.go index d31b0f8..31b634d 100644 --- a/gor.go +++ b/gor.go @@ -80,9 +80,9 @@ func main() { profileCPU(*cpuprofile) } - if Settings.Pprof != "" { + if Settings.pprof() != "" { go func() { - log.Println(http.ListenAndServe(Settings.Pprof, nil)) + log.Println(http.ListenAndServe(Settings.pprof(), nil)) }() } @@ -95,16 +95,16 @@ func main() { os.Exit(1) }() - if Settings.ExitAfter > 0 { - log.Println("Running gor for a duration of", Settings.ExitAfter) + if Settings.exitAfter() > 0 { + log.Println("Running gor for a duration of", Settings.exitAfter()) - time.AfterFunc(Settings.ExitAfter, func() { - log.Println("Stopping gor after", Settings.ExitAfter) + time.AfterFunc(Settings.exitAfter(), func() { + log.Println("Stopping gor after", Settings.exitAfter()) close(closeCh) }) } - emitter.Start(plugins, Settings.Middleware) + emitter.Start(plugins, Settings.middleware()) } func finalize(plugins *InOutPlugins) { diff --git a/gor_stat.go b/gor_stat.go index a761f19..85fbea7 100644 --- a/gor_stat.go +++ b/gor_stat.go @@ -25,7 +25,7 @@ func NewGorStat(statName string, rateMs int) (s *GorStat) { s.max = 0 s.count = 0 - if Settings.Stats { + if Settings.stats() { log.Println(s.statName + ":latest,mean,max,count,count/second,gcount") go s.reportStats() } @@ -33,7 +33,7 @@ func NewGorStat(statName string, rateMs int) (s *GorStat) { } func (s *GorStat) Write(latest int) { - if Settings.Stats { + if Settings.stats() { if latest > s.max { s.max = latest } diff --git a/input_file_test.go b/input_file_test.go index 7858ff2..32d66b8 100644 --- a/input_file_test.go +++ b/input_file_test.go @@ -329,7 +329,7 @@ func CreateCaptureFile(requestGenerator *RequestGenerator) *CaptureFile { plugins.All = append(plugins.All, output, outputFile) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) requestGenerator.emit() requestGenerator.wg.Wait() @@ -359,7 +359,7 @@ func ReadFromCaptureFile(captureFile *os.File, count int, callback writeCallback wg.Add(count) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) done := make(chan int, 1) go func() { diff --git a/input_http_test.go b/input_http_test.go index f0cad6c..9c62302 100644 --- a/input_http_test.go +++ b/input_http_test.go @@ -28,7 +28,7 @@ func TestHTTPInput(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) address := strings.Replace(input.listener.Addr().String(), "[::]", "127.0.0.1", -1) @@ -65,7 +65,7 @@ func TestInputHTTPLargePayload(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) wg.Add(1) address := strings.Replace(input.listener.Addr().String(), "[::]", "127.0.0.1", -1) diff --git a/input_kafka.go b/input_kafka.go index 2194439..fa0eb0e 100644 --- a/input_kafka.go +++ b/input_kafka.go @@ -61,7 +61,7 @@ func NewKafkaInput(address string, config *InputKafkaConfig) *KafkaInput { } }(consumer) - if Settings.Verbose { + if Settings.verbose() { // Start infinite loop for tracking errors for kafka producer. go i.ErrorHandler(consumer) } diff --git a/input_raw.go b/input_raw.go index d29f49a..1323899 100644 --- a/input_raw.go +++ b/input_raw.go @@ -98,7 +98,7 @@ func (i *RAWInput) listen(address string) { log.Fatalf("input-raw: error while parsing address: %s", err) } - i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.protocol, i.bpfFilter, i.timestampType, i.bufferSize, Settings.InputRAWConfig.OverrideSnapLen, Settings.InputRAWConfig.ImmediateMode) + i.listener = raw.NewListener(host, port, i.engine, i.trackResponse, i.expire, i.protocol, i.bpfFilter, i.timestampType, i.bufferSize, Settings.inputRAWConfigOverrideSnapLen(), Settings.inputRAWConfigImmediateMode()) ch := i.listener.Receiver() diff --git a/input_raw_test.go b/input_raw_test.go index 314997d..44ed2ee 100644 --- a/input_raw_test.go +++ b/input_raw_test.go @@ -58,7 +58,7 @@ func TestRAWInputIPv4(t *testing.T) { atomic.AddInt64(&respCounter, 1) } - if Settings.Debug { + if Settings.debug() { log.Println(reqCounter, respCounter) } @@ -74,7 +74,7 @@ func TestRAWInputIPv4(t *testing.T) { client := NewHTTPClient("http://"+listener.Addr().String(), &HTTPClientConfig{}) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { // request + response @@ -125,7 +125,7 @@ func TestRAWInputNoKeepAlive(t *testing.T) { client := NewHTTPClient("http://"+listener.Addr().String(), &HTTPClientConfig{}) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { // request + response @@ -168,7 +168,7 @@ func TestRAWInputIPv6(t *testing.T) { atomic.AddInt64(&respCounter, 1) } - if Settings.Debug { + if Settings.debug() { log.Println(reqCounter, respCounter) } @@ -184,7 +184,7 @@ func TestRAWInputIPv6(t *testing.T) { client := NewHTTPClient("http://"+listener.Addr().String(), &HTTPClientConfig{}) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { // request + response @@ -251,7 +251,7 @@ func TestInputRAW100Expect(t *testing.T) { plugins.All = append(plugins.All, input, testOutput, httpOutput) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) // Origin + Response/Request Test Output + Request Http Output wg.Add(4) @@ -305,7 +305,7 @@ func TestInputRAWChunkedEncoding(t *testing.T) { plugins.All = append(plugins.All, input, httpOutput) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) wg.Add(2) curl := exec.Command("curl", "http://"+originAddr, "--header", "Transfer-Encoding: chunked", "--header", "Expect:", "--data-binary", "@README.md") @@ -373,7 +373,7 @@ func TestInputRAWLargePayload(t *testing.T) { plugins.All = append(plugins.All, input, httpOutput) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) wg.Add(2) curl := exec.Command("curl", "http://"+originAddr, "--header", "Transfer-Encoding: chunked", "--header", "Expect:", "--data-binary", "@/tmp/large") @@ -421,7 +421,7 @@ func BenchmarkRAWInput(b *testing.B) { plugins.All = append(plugins.All, input, output, httpOutput) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) emitted := 0 fileContent, _ := ioutil.ReadFile("LICENSE.txt") diff --git a/input_tcp_test.go b/input_tcp_test.go index f470086..7b892da 100644 --- a/input_tcp_test.go +++ b/input_tcp_test.go @@ -34,7 +34,7 @@ func TestTCPInput(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) tcpAddr, err := net.ResolveTCPAddr("tcp", input.listener.Addr().String()) @@ -117,7 +117,7 @@ func TestTCPInputSecure(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) conf := &tls.Config{ InsecureSkipVerify: true, diff --git a/limiter_test.go b/limiter_test.go index 90d0174..29d3760 100644 --- a/limiter_test.go +++ b/limiter_test.go @@ -25,7 +25,7 @@ func TestOutputLimiter(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { input.EmitGET() @@ -52,7 +52,7 @@ func TestInputLimiter(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { input.(*Limiter).plugin.(*TestInput).EmitGET() @@ -79,7 +79,7 @@ func TestPercentLimiter1(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { input.EmitGET() @@ -107,7 +107,7 @@ func TestPercentLimiter2(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { input.EmitGET() diff --git a/middleware.go b/middleware.go index 000120b..95e4b24 100644 --- a/middleware.go +++ b/middleware.go @@ -75,7 +75,7 @@ func (m *Middleware) copy(to io.Writer, from io.Reader) { payload := buf[0:nr] - if Settings.PrettifyHTTP { + if Settings.prettifyHTTP() { payload = prettifyHTTP(payload) nr = len(payload) @@ -84,7 +84,7 @@ func (m *Middleware) copy(to io.Writer, from io.Reader) { } } - if Settings.PrettifyHTTP { + if Settings.prettifyHTTP() { payload = prettifyHTTP(payload) nr = len(payload) } @@ -96,7 +96,7 @@ func (m *Middleware) copy(to io.Writer, from io.Reader) { to.Write(dst[0 : nr*2+1]) m.mu.Unlock() - if Settings.Debug { + if Settings.debug() { Debug("[MIDDLEWARE-MASTER] Sending:", string(buf[0:nr]), "From:", from) } } @@ -121,7 +121,7 @@ func (m *Middleware) read(from io.Reader) { fmt.Fprintln(os.Stderr, "Failed to decode input payload", err, len(line), string(line[:len(line)-1])) } - if Settings.Debug { + if Settings.debug() { Debug("[MIDDLEWARE-MASTER] Received:", string(buf)) } diff --git a/middleware_test.go b/middleware_test.go index a7c4781..620fbc3 100644 --- a/middleware_test.go +++ b/middleware_test.go @@ -114,7 +114,7 @@ func TestEchoMiddleware(t *testing.T) { quit := make(chan int) - Settings.Middleware = "./examples/middleware/echo.sh" + Settings.setMiddleware("./examples/middleware/echo.sh") // Catch traffic from one service fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) @@ -132,7 +132,7 @@ func TestEchoMiddleware(t *testing.T) { // Start Gor emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) // Wait till middleware initialization time.Sleep(100 * time.Millisecond) @@ -153,7 +153,7 @@ func TestEchoMiddleware(t *testing.T) { emitter.Close() time.Sleep(200 * time.Millisecond) - Settings.Middleware = "" + Settings.setMiddleware("") } func TestTokenMiddleware(t *testing.T) { @@ -180,7 +180,7 @@ func TestTokenMiddleware(t *testing.T) { quit := make(chan int) - Settings.Middleware = "go run ./examples/middleware/token_modifier.go" + Settings.setMiddleware("go run ./examples/middleware/token_modifier.go") fromAddr := strings.Replace(from.Listener.Addr().String(), "[::]", "127.0.0.1", -1) // Catch traffic from one service @@ -198,7 +198,7 @@ func TestTokenMiddleware(t *testing.T) { // Start Gor emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) // Wait for middleware to initialize // Give go compiller time to build programm @@ -225,5 +225,5 @@ func TestTokenMiddleware(t *testing.T) { wg.Wait() emitter.Close() time.Sleep(100 * time.Millisecond) - Settings.Middleware = "" + Settings.setMiddleware("") } diff --git a/output_file_test.go b/output_file_test.go index 285da28..927f123 100644 --- a/output_file_test.go +++ b/output_file_test.go @@ -27,7 +27,7 @@ func TestFileOutput(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { wg.Add(2) @@ -53,7 +53,7 @@ func TestFileOutput(t *testing.T) { quit2 := make(chan int) emitter2 := NewEmitter(quit2) - go emitter2.Start(plugins2, Settings.Middleware) + go emitter2.Start(plugins2, Settings.middleware()) wg.Wait() emitter2.Close() diff --git a/output_http.go b/output_http.go index d557981..6574153 100644 --- a/output_http.go +++ b/output_http.go @@ -143,7 +143,7 @@ func NewHTTPOutput(address string, config *HTTPOutputConfig) io.Writer { o.elasticSearch.Init(o.config.ElasticSearch) } - if Settings.RecognizeTCPSessions { + if Settings.recognizeTCPSessions() { if !PRO { log.Fatal("Detailed TCP sessions work only with PRO license") } @@ -249,7 +249,7 @@ func (o *HTTPOutput) Write(data []byte) (n int, err error) { o.queueStats.Write(len(o.queue)) } - if !Settings.RecognizeTCPSessions && o.config.WorkersMax != o.config.WorkersMin { + if !Settings.recognizeTCPSessions() && o.config.WorkersMax != o.config.WorkersMin { workersCount := int(atomic.LoadInt64(&o.activeWorkers)) if len(o.queue) > workersCount { @@ -275,7 +275,7 @@ func (o *HTTPOutput) Read(data []byte) (int, error) { case resp = <-o.responses: } - if Settings.Debug { + if Settings.debug() { Debug("[OUTPUT-HTTP] Received response:", string(resp.payload)) } @@ -289,7 +289,7 @@ func (o *HTTPOutput) Read(data []byte) (int, error) { func (o *HTTPOutput) sendRequest(client *HTTPClient, request []byte) { meta := payloadMeta(request) - if Settings.Debug { + if Settings.debug() { Debug(meta) } diff --git a/output_http_test.go b/output_http_test.go index bb2260a..1b82c88 100644 --- a/output_http_test.go +++ b/output_http_test.go @@ -42,7 +42,7 @@ func TestHTTPOutput(t *testing.T) { headers := HTTPHeaders{HTTPHeader{"User-Agent", "Gor"}} methods := HTTPMethods{[]byte("GET"), []byte("PUT"), []byte("POST")} - Settings.ModifierConfig = HTTPModifierConfig{Headers: headers, Methods: methods} + Settings.setModifierConfig(HTTPModifierConfig{Headers: headers, Methods: methods}) http_output := NewHTTPOutput(server.URL, &HTTPOutputConfig{Debug: true, TrackResponses: true}) output := NewTestOutput(func(data []byte) { @@ -56,7 +56,7 @@ func TestHTTPOutput(t *testing.T) { plugins.All = append(plugins.All, input, output, http_output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 10; i++ { // 2 http-output, 2 - test output request, 2 - test output http response @@ -75,7 +75,7 @@ func TestHTTPOutput(t *testing.T) { t.Error("Should create workers for each request", activeWorkers) } - Settings.ModifierConfig = HTTPModifierConfig{} + Settings.setModifierConfig(HTTPModifierConfig{}) } func TestHTTPOutputKeepOriginalHost(t *testing.T) { @@ -94,7 +94,7 @@ func TestHTTPOutputKeepOriginalHost(t *testing.T) { defer server.Close() headers := HTTPHeaders{HTTPHeader{"Host", "custom-host.com"}} - Settings.ModifierConfig = HTTPModifierConfig{Headers: headers} + Settings.setModifierConfig(HTTPModifierConfig{Headers: headers}) output := NewHTTPOutput(server.URL, &HTTPOutputConfig{Debug: false, OriginalHost: true}) @@ -105,14 +105,14 @@ func TestHTTPOutputKeepOriginalHost(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) wg.Add(1) input.EmitGET() wg.Wait() emitter.Close() - Settings.ModifierConfig = HTTPModifierConfig{} + Settings.setModifierConfig(HTTPModifierConfig{}) } func TestHTTPOutputSSL(t *testing.T) { @@ -134,7 +134,7 @@ func TestHTTPOutputSSL(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) wg.Add(2) @@ -157,7 +157,7 @@ func TestHTTPOutputSessions(t *testing.T) { })) defer server.Close() - Settings.RecognizeTCPSessions = true + Settings.setRecognizeTCPSessions(true) output := NewHTTPOutput(server.URL, &HTTPOutputConfig{Debug: true}) plugins := &InOutPlugins{ @@ -165,7 +165,7 @@ func TestHTTPOutputSessions(t *testing.T) { Outputs: []io.Writer{output}, } emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) uuid1 := []byte("1234567890123456789a0000") uuid2 := []byte("1234567890123456789d0000") @@ -190,7 +190,7 @@ func TestHTTPOutputSessions(t *testing.T) { emitter.Close() - Settings.RecognizeTCPSessions = false + Settings.setRecognizeTCPSessions(false) } func BenchmarkHTTPOutput(b *testing.B) { @@ -213,7 +213,7 @@ func BenchmarkHTTPOutput(b *testing.B) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < b.N; i++ { wg.Add(1) diff --git a/output_kafka.go b/output_kafka.go index 6616a88..4f27d76 100644 --- a/output_kafka.go +++ b/output_kafka.go @@ -49,7 +49,7 @@ func NewKafkaOutput(address string, config *OutputKafkaConfig) io.Writer { producer: producer, } - if Settings.Verbose { + if Settings.verbose() { // Start infinite loop for tracking errors for kafka producer. go o.ErrorHandler() } diff --git a/output_tcp.go b/output_tcp.go index f283d70..9b943fe 100644 --- a/output_tcp.go +++ b/output_tcp.go @@ -34,7 +34,7 @@ func NewTCPOutput(address string, config *TCPOutputConfig) io.Writer { o.address = address o.config = config - if Settings.OutputTCPStats { + if Settings.outputTCPStats() { o.bufStats = NewGorStat("output_tcp", 5000) } @@ -114,7 +114,7 @@ func (o *TCPOutput) Write(data []byte) (n int, err error) { bufferIndex := o.getBufferIndex(data) o.buf[bufferIndex] <- newBuf - if Settings.OutputTCPStats { + if Settings.outputTCPStats() { o.bufStats.Write(len(o.buf[bufferIndex])) } diff --git a/output_tcp_test.go b/output_tcp_test.go index d60d20c..d0bdb41 100644 --- a/output_tcp_test.go +++ b/output_tcp_test.go @@ -27,7 +27,7 @@ func TestTCPOutput(t *testing.T) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) for i := 0; i < 100; i++ { wg.Add(1) @@ -82,7 +82,7 @@ func BenchmarkTCPOutput(b *testing.B) { plugins.All = append(plugins.All, input, output) emitter := NewEmitter(quit) - go emitter.Start(plugins, Settings.Middleware) + go emitter.Start(plugins, Settings.middleware()) b.ResetTimer() for i := 0; i < b.N; i++ { diff --git a/plugins.go b/plugins.go index 1b631a6..586572f 100644 --- a/plugins.go +++ b/plugins.go @@ -82,46 +82,46 @@ func InitPlugins() *InOutPlugins { pluginMu.Lock() defer pluginMu.Unlock() - for _, options := range Settings.InputDummy { + for _, options := range Settings.inputDummy() { registerPlugin(NewDummyInput, options) } - for range Settings.OutputDummy { + for range Settings.outputDummy() { registerPlugin(NewDummyOutput) } - if Settings.OutputStdout { + if Settings.outputStdout() { registerPlugin(NewDummyOutput) } - if Settings.OutputNull { + if Settings.outputNull() { registerPlugin(NewNullOutput) } engine := EnginePcap - if Settings.InputRAWConfig.Engine == "raw_socket" { + if Settings.inputRAWConfigEngine() == "raw_socket" { engine = EngineRawSocket - } else if Settings.InputRAWConfig.Engine == "pcap_file" { + } else if Settings.inputRAWConfigEngine() == "pcap_file" { engine = EnginePcapFile } - for _, options := range Settings.InputRAW { - registerPlugin(NewRAWInput, options, engine, Settings.InputRAWConfig.TrackResponse, Settings.InputRAWConfig.Expire, Settings.InputRAWConfig.RealIPHeader, Settings.InputRAWConfig.Protocol, Settings.InputRAWConfig.BpfFilter, Settings.InputRAWConfig.TimestampType, Settings.InputRAWConfig.bufferSize) + for _, options := range Settings.inputRAW() { + registerPlugin(NewRAWInput, options, engine, Settings.inputRAWConfigTrackResponse(), Settings.inputRAWConfigExpire, Settings.inputRAWConfigRealIPHeader, Settings.inputRAWConfigProtocol, Settings.inputRAWConfigBpfFilter, Settings.inputRAWConfigTimestampType, Settings.InputRAWConfig.bufferSize) } - for _, options := range Settings.InputTCP { + for _, options := range Settings.inputTCP() { registerPlugin(NewTCPInput, options, &Settings.InputTCPConfig) } - for _, options := range Settings.OutputTCP { + for _, options := range Settings.outputTCP() { registerPlugin(NewTCPOutput, options, &Settings.OutputTCPConfig) } - for _, options := range Settings.InputFile { - registerPlugin(NewFileInput, options, Settings.InputFileLoop) + for _, options := range Settings.inputFile() { + registerPlugin(NewFileInput, options, Settings.inputFileLoop()) } - for _, path := range Settings.OutputFile { + for _, path := range Settings.outputFile() { if strings.HasPrefix(path, "s3://") { registerPlugin(NewS3Output, path, &Settings.OutputFileConfig) } else { diff --git a/plugins_test.go b/plugins_test.go index 0cfbb97..a4cc57e 100644 --- a/plugins_test.go +++ b/plugins_test.go @@ -5,10 +5,10 @@ import ( ) func TestPluginsRegistration(t *testing.T) { - Settings.InputDummy = MultiOption{"[]"} - Settings.OutputDummy = MultiOption{"[]"} - Settings.OutputHTTP = MultiOption{"www.example.com|10"} - Settings.InputFile = MultiOption{"/dev/null"} + Settings.setInputDummy(MultiOption{"[]"}) + Settings.setOutputDummy(MultiOption{"[]"}) + Settings.setOutputHTTP(MultiOption{"www.example.com|10"}) + Settings.setInputFile(MultiOption{"/dev/null"}) plugins := InitPlugins() diff --git a/settings.go b/settings.go index 431d664..c5b7c38 100644 --- a/settings.go +++ b/settings.go @@ -107,6 +107,1052 @@ type AppSettings struct { ConfigFile string `json:"config-file" mapstructure:"config-file"` ConfigServerAddress string `json:"config-server-address" mapstructure:"config-server-address"` RemoteConfigHost string `json:"config-remote-host" mapstructure:"config-remote-host"` + + sync.RWMutex +} + +func (a *AppSettings) verbose() bool { + a.RLock() + defer a.RUnlock() + return a.Verbose +} + +func (a *AppSettings) setVerbose(verbose bool) { + a.Lock() + defer a.Unlock() + a.Verbose = verbose +} + +func (a *AppSettings) debug() bool { + a.RLock() + defer a.RUnlock() + return a.Debug +} + +func (a *AppSettings) setDebug(debug bool) { + a.Lock() + defer a.Unlock() + a.Debug = debug +} + +func (a *AppSettings) stats() bool { + a.RLock() + defer a.RUnlock() + return a.Stats +} + +func (a *AppSettings) setStats(stats bool) { + a.Lock() + defer a.Unlock() + a.Stats = stats +} + +func (a *AppSettings) exitAfter() time.Duration { + a.RLock() + defer a.RUnlock() + return a.ExitAfter +} + +func (a *AppSettings) setExitAfter(exitAfter time.Duration) { + a.Lock() + defer a.Unlock() + a.ExitAfter = exitAfter +} + +func (a *AppSettings) splitOutput() bool { + a.RLock() + defer a.RUnlock() + return a.SplitOutput +} + +func (a *AppSettings) setSplitOutput(splitOutput bool) { + a.Lock() + defer a.Unlock() + a.SplitOutput = splitOutput +} + +func (a *AppSettings) recognizeTCPSessions() bool { + a.RLock() + defer a.RUnlock() + return a.RecognizeTCPSessions +} + +func (a *AppSettings) setRecognizeTCPSessions(recognizeTCPSessions bool) { + a.Lock() + defer a.Unlock() + a.RecognizeTCPSessions = recognizeTCPSessions +} + +func (a *AppSettings) pprof() string { + a.RLock() + defer a.RUnlock() + return a.Pprof +} + +func (a *AppSettings) setPprof(pprof string) { + a.Lock() + defer a.Unlock() + a.Pprof = pprof +} + +func (a *AppSettings) inputDummy() MultiOption { + a.RLock() + defer a.RUnlock() + return a.InputDummy +} + +func (a *AppSettings) setInputDummy(inputDummy MultiOption) { + a.Lock() + defer a.Unlock() + a.InputDummy = inputDummy +} + +func (a *AppSettings) outputDummy() MultiOption { + a.RLock() + defer a.RUnlock() + return a.OutputDummy +} + +func (a *AppSettings) setOutputDummy(outputDummy MultiOption) { + a.Lock() + defer a.Unlock() + a.OutputDummy = outputDummy +} + +func (a *AppSettings) outputStdout() bool { + a.RLock() + defer a.RUnlock() + return a.OutputStdout +} + +func (a *AppSettings) setOutputStdout(outputStdout bool) { + a.Lock() + defer a.Unlock() + a.OutputStdout = outputStdout +} + +func (a *AppSettings) outputNull() bool { + a.RLock() + defer a.RUnlock() + return a.OutputNull +} + +func (a *AppSettings) setOutputNull(outputNull bool) { + a.Lock() + defer a.Unlock() + a.OutputNull = outputNull +} + +func (a *AppSettings) inputTCP() MultiOption { + a.RLock() + defer a.RUnlock() + return a.InputTCP +} + +func (a *AppSettings) setInputTCP(inputTCP MultiOption) { + a.Lock() + defer a.Unlock() + a.InputTCP = inputTCP +} + +func (a *AppSettings) inputTCPConfig() TCPInputConfig { + a.RLock() + defer a.RUnlock() + return a.InputTCPConfig +} + +func (a *AppSettings) setInputTCPConfig(inputTCPConfig TCPInputConfig) { + a.Lock() + defer a.Unlock() + a.InputTCPConfig = inputTCPConfig +} + +func (a *AppSettings) inputTCPConfigSecure() bool { + a.RLock() + defer a.RUnlock() + return a.InputTCPConfig.Secure +} + +func (a *AppSettings) setInputTCPConfigSecure(secure bool) { + a.Lock() + defer a.Unlock() + a.InputTCPConfig.Secure = secure +} + +func (a *AppSettings) inputTCPConfigCertificatePath() string { + a.RLock() + defer a.RUnlock() + return a.InputTCPConfig.CertificatePath +} + +func (a *AppSettings) setInputTCPConfigCertificatePath(certificatePath string) { + a.Lock() + defer a.Unlock() + a.InputTCPConfig.CertificatePath = certificatePath +} + +func (a *AppSettings) inputTCPConfigKeyPath() string { + a.RLock() + defer a.RUnlock() + return a.InputTCPConfig.KeyPath +} + +func (a *AppSettings) setInputTCPConfigKeyPath(keyPath string) { + a.Lock() + defer a.Unlock() + a.InputTCPConfig.KeyPath = keyPath +} + +func (a *AppSettings) outputTCP() MultiOption { + a.RLock() + defer a.RUnlock() + return a.OutputTCP +} + +func (a *AppSettings) setOutputTCP(outputTCP MultiOption) { + a.Lock() + defer a.Unlock() + a.OutputTCP = outputTCP +} + +func (a *AppSettings) outputTCPConfig() TCPOutputConfig { + a.RLock() + defer a.RUnlock() + return a.OutputTCPConfig +} + +func (a *AppSettings) setOutputTCPConfig(outputTCPConfig TCPOutputConfig) { + a.Lock() + defer a.Unlock() + a.OutputTCPConfig = outputTCPConfig +} + +func (a *AppSettings) outputTCPConfigSecure() bool { + a.RLock() + defer a.RUnlock() + return a.OutputTCPConfig.Secure +} + +func (a *AppSettings) setOutputTCPConfigSecure(secure bool) { + a.Lock() + defer a.Unlock() + a.OutputTCPConfig.Secure = secure +} + +func (a *AppSettings) outputTCPConfigSticky() bool { + a.RLock() + defer a.RUnlock() + return a.OutputTCPConfig.Sticky +} + +func (a *AppSettings) setOutputTCPConfigSticky(sticky bool) { + a.Lock() + defer a.Unlock() + a.OutputTCPConfig.Sticky = sticky +} + +func (a *AppSettings) outputTCPStats() bool { + a.RLock() + defer a.RUnlock() + return a.OutputTCPStats +} + +func (a *AppSettings) setOutputTCPStats(outputTCPStats bool) { + a.Lock() + defer a.Unlock() + a.OutputTCPStats = outputTCPStats +} + +func (a *AppSettings) inputFile() MultiOption { + a.RLock() + defer a.RUnlock() + return a.InputFile +} + +func (a *AppSettings) setInputFile(inputFile MultiOption) { + a.Lock() + defer a.Unlock() + a.InputFile = inputFile +} + +func (a *AppSettings) inputFileLoop() bool { + a.RLock() + defer a.RUnlock() + return a.InputFileLoop +} + +func (a *AppSettings) setInputFileLoop(inputFileLoop bool) { + a.Lock() + defer a.Unlock() + a.InputFileLoop = inputFileLoop +} + +func (a *AppSettings) outputFile() MultiOption { + a.RLock() + defer a.RUnlock() + return a.OutputFile +} + +func (a *AppSettings) setOutputFile(outputFile MultiOption) { + a.Lock() + defer a.Unlock() + a.OutputFile = outputFile +} + +func (a *AppSettings) outputFileConfigFlushInterval() time.Duration { + a.RLock() + defer a.RUnlock() + return a.OutputFileConfig.FlushInterval +} + +func (a *AppSettings) setOutputFileConfigFlushInterval(flushInterval time.Duration) { + a.Lock() + defer a.Unlock() + a.OutputFileConfig.FlushInterval = flushInterval +} + +func (a *AppSettings) outputFileConfigQueueLimit() int64 { + a.RLock() + defer a.RUnlock() + return a.OutputFileConfig.QueueLimit +} + +func (a *AppSettings) setOutputFileConfigQueueLimit(queueLimit int64) { + a.Lock() + defer a.Unlock() + a.OutputFileConfig.QueueLimit = queueLimit +} + +func (a *AppSettings) outputFileConfigAppend() bool { + a.RLock() + defer a.RUnlock() + return a.OutputFileConfig.Append +} + +func (a *AppSettings) setOutputFileConfigAppend(append bool) { + a.Lock() + defer a.Unlock() + a.OutputFileConfig.Append = append +} + +func (a *AppSettings) outputFileConfigBufferPath() string { + a.RLock() + defer a.RUnlock() + return a.OutputFileConfig.BufferPath +} + +func (a *AppSettings) setOutputFileConfigBufferPath(bufferPath string) { + a.Lock() + defer a.Unlock() + a.OutputFileConfig.BufferPath = bufferPath +} + +func (a *AppSettings) inputRAW() MultiOption { + a.RLock() + defer a.RUnlock() + return a.InputRAW +} + +func (a *AppSettings) setInputRAW(inputRAW MultiOption) { + a.Lock() + defer a.Unlock() + a.InputRAW = inputRAW +} + +func (a *AppSettings) inputRAWConfigEngine() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.Engine +} + +func (a *AppSettings) setInputRAWConfigEngine(engine string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.Engine = engine +} + +func (a *AppSettings) inputRAWConfigTrackResponse() bool { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.TrackResponse +} + +func (a *AppSettings) setInputRAWConfigTrackResponse(trackResponse bool) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.TrackResponse = trackResponse +} + +func (a *AppSettings) inputRAWConfigRealIPHeader() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.RealIPHeader +} + +func (a *AppSettings) setInputRAWConfigRealIPHeader(realIPHeader string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.RealIPHeader = realIPHeader +} + +func (a *AppSettings) inputRAWConfigExpire() time.Duration { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.Expire +} + +func (a *AppSettings) setInputRAWConfigExpire(expire time.Duration) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.Expire = expire +} + +func (a *AppSettings) inputRAWConfigProtocol() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.Protocol +} + +func (a *AppSettings) setInputRAWConfigProtocol(protocol string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.Protocol = protocol +} + +func (a *AppSettings) inputRAWConfigBpfFilter() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.BpfFilter +} + +func (a *AppSettings) setInputRAWConfigBpfFilter(bpfFilter string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.BpfFilter = bpfFilter +} + +func (a *AppSettings) inputRAWConfigTimestampType() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.TimestampType +} + +func (a *AppSettings) setInputRAWConfigTimestampType(timestampType string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.TimestampType = timestampType +} + +func (a *AppSettings) inputRAWConfigImmediateMode() bool { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.ImmediateMode +} + +func (a *AppSettings) setInputRAWConfigImmediateMode(immediateMode bool) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.ImmediateMode = immediateMode +} + +func (a *AppSettings) inputRAWConfigOverrideSnapLen() bool { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.OverrideSnapLen +} + +func (a *AppSettings) setInputRAWConfigOverrideSnapLen(overrideSnapLen bool) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.OverrideSnapLen = overrideSnapLen +} + +func (a *AppSettings) inputRAWConfigBufferSizeFlag() string { + a.RLock() + defer a.RUnlock() + return a.InputRAWConfig.BufferSizeFlag +} + +func (a *AppSettings) setInputRAWConfigBufferSizeFlag(bufferSizeFlag string) { + a.Lock() + defer a.Unlock() + a.InputRAWConfig.BufferSizeFlag = bufferSizeFlag +} + +func (a *AppSettings) outputFileSizeFlag() string { + a.RLock() + defer a.RUnlock() + return a.OutputFileSizeFlag +} + +func (a *AppSettings) setOutputFileSizeFlag(outputFileSizeFlag string) { + a.Lock() + defer a.Unlock() + a.OutputFileSizeFlag = outputFileSizeFlag +} + +func (a *AppSettings) outputFileMaxSizeFlag() string { + a.RLock() + defer a.RUnlock() + return a.OutputFileMaxSizeFlag +} + +func (a *AppSettings) setOutputFileMaxSizeFlag(outputFileMaxSizeFlag string) { + a.Lock() + defer a.Unlock() + a.OutputFileMaxSizeFlag = outputFileMaxSizeFlag +} + +func (a *AppSettings) copyBufferSizeFlag() string { + a.RLock() + defer a.RUnlock() + return a.CopyBufferSizeFlag +} + +func (a *AppSettings) setCopyBufferSizeFlag(copyBufferSizeFlag string) { + a.Lock() + defer a.Unlock() + a.CopyBufferSizeFlag = copyBufferSizeFlag +} + +func (a *AppSettings) middleware() string { + a.RLock() + defer a.RUnlock() + return a.Middleware +} + +func (a *AppSettings) setMiddleware(middleware string) { + a.Lock() + defer a.Unlock() + a.Middleware = middleware +} + +func (a *AppSettings) inputHTTP() MultiOption { + a.RLock() + defer a.RUnlock() + return a.InputHTTP +} + +func (a *AppSettings) setInputHTTP(inputHTTP MultiOption) { + a.Lock() + defer a.Unlock() + a.InputHTTP = inputHTTP +} + +func (a *AppSettings) outputHTTP() MultiOption { + a.RLock() + defer a.RUnlock() + return a.OutputHTTP +} + +func (a *AppSettings) setOutputHTTP(outputHTTP MultiOption) { + a.Lock() + defer a.Unlock() + a.OutputHTTP = outputHTTP +} + +func (a *AppSettings) prettifyHTTP() bool { + a.RLock() + defer a.RUnlock() + return a.PrettifyHTTP +} + +func (a *AppSettings) setPrettifyHTTP(prettifyHTTP bool) { + a.Lock() + defer a.Unlock() + a.PrettifyHTTP = prettifyHTTP +} + +func (a *AppSettings) outputHTTPConfigRedirectLimit() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.RedirectLimit +} + +func (a *AppSettings) setOutputHTTPConfigRedirectLimit(redirectLimit int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.RedirectLimit = redirectLimit +} + +func (a *AppSettings) outputHTTPConfigStats() bool { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.Stats +} + +func (a *AppSettings) setOutputHTTPConfigStats(stats bool) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.Stats = stats +} + +func (a *AppSettings) outputHTTPConfigWorkersMin() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.WorkersMin +} + +func (a *AppSettings) setOutputHTTPConfigWorkersMin(workersMin int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.WorkersMin = workersMin +} + +func (a *AppSettings) outputHTTPConfigWorkersMax() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.WorkersMax +} + +func (a *AppSettings) setOutputHTTPConfigWorkersMax(workersMax int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.WorkersMax = workersMax +} + +func (a *AppSettings) outputHTTPConfigStatsMs() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.StatsMs +} + +func (a *AppSettings) setOutputHTTPConfigStatsMs(statsMs int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.StatsMs = statsMs +} + +func (a *AppSettings) outputHTTPConfigQueueLen() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.QueueLen +} + +func (a *AppSettings) setOutputHTTPConfigQueueLen(queueLen int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.QueueLen = queueLen +} + +func (a *AppSettings) outputHTTPConfigElasticSearch() string { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.ElasticSearch +} + +func (a *AppSettings) setOutputHTTPConfigElasticSearch(elasticSearch string) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.ElasticSearch = elasticSearch +} + +func (a *AppSettings) outputHTTPConfigTimeout() time.Duration { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.Timeout +} + +func (a *AppSettings) setOutputHTTPConfigTimeout(timeout time.Duration) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.Timeout = timeout +} + +func (a *AppSettings) outputHTTPConfigOriginalHost() bool { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.OriginalHost +} + +func (a *AppSettings) setOutputHTTPConfigOriginalHost(host bool) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.OriginalHost = host +} + +func (a *AppSettings) outputHTTPConfigBufferSize() int { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.BufferSize +} + +func (a *AppSettings) setOutputHTTPConfigBufferSize(bufferSize int) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.BufferSize = bufferSize +} + +func (a *AppSettings) outputHTTPConfigCompatibilityMode() bool { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.CompatibilityMode +} + +func (a *AppSettings) setOutputHTTPConfigCompatibilityMode(compatibilityMode bool) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.CompatibilityMode = compatibilityMode +} + +func (a *AppSettings) outputHTTPConfigDebug() bool { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.Debug +} + +func (a *AppSettings) setOutputHTTPConfigDebug(debug bool) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.Debug = debug +} + +func (a *AppSettings) outputHTTPConfigTrackResponses() bool { + a.RLock() + defer a.RUnlock() + return a.OutputHTTPConfig.TrackResponses +} + +func (a *AppSettings) setOutputHTTPConfigTrackResponses(trackResponses bool) { + a.Lock() + defer a.Unlock() + a.OutputHTTPConfig.TrackResponses = trackResponses +} + +func (a *AppSettings) outputBinary() MultiOption { + a.RLock() + defer a.RUnlock() + return a.OutputBinary +} + +func (a *AppSettings) setOutputBinary(outputBinary MultiOption) { + a.Lock() + defer a.Unlock() + a.OutputBinary = outputBinary +} + +func (a *AppSettings) outputBinaryConfigWorkers() int { + a.RLock() + defer a.RUnlock() + return a.OutputBinaryConfig.Workers +} + +func (a *AppSettings) setOutputBinaryConfigWorkers(workers int) { + a.Lock() + defer a.Unlock() + a.OutputBinaryConfig.Workers = workers +} + +func (a *AppSettings) outputBinaryConfigTimeout() time.Duration { + a.RLock() + defer a.RUnlock() + return a.OutputBinaryConfig.Timeout +} + +func (a *AppSettings) setOutputBinaryConfigTimeout(timeout time.Duration) { + a.Lock() + defer a.Unlock() + a.OutputBinaryConfig.Timeout = timeout +} + +func (a *AppSettings) outputBinaryConfigBufferSize() int { + a.RLock() + defer a.RUnlock() + return a.OutputBinaryConfig.BufferSize +} + +func (a *AppSettings) setOutputBinaryConfigBufferSize(bufferSize int) { + a.Lock() + defer a.Unlock() + a.OutputBinaryConfig.BufferSize = bufferSize +} + +func (a *AppSettings) outputBinaryConfigDebug() bool { + a.RLock() + defer a.RUnlock() + return a.OutputBinaryConfig.Debug +} + +func (a *AppSettings) setOutputBinaryConfigDebug(debug bool) { + a.Lock() + defer a.Unlock() + a.OutputBinaryConfig.Debug = debug +} + +func (a *AppSettings) outputBinaryConfigTrackResponses() bool { + a.RLock() + defer a.RUnlock() + return a.OutputBinaryConfig.TrackResponses +} + +func (a *AppSettings) setOutputBinaryConfigTrackResponses(trackResponses bool) { + a.Lock() + defer a.Unlock() + a.OutputBinaryConfig.TrackResponses = trackResponses +} + +func (a *AppSettings) modifierConfig() HTTPModifierConfig { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig +} + +func (a *AppSettings) setModifierConfig(config HTTPModifierConfig) { + a.Lock() + defer a.Unlock() + a.ModifierConfig = config +} + +func (a *AppSettings) modifierConfigUrlNegativeRegexp() HTTPUrlRegexp { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.UrlNegativeRegexp +} + +func (a *AppSettings) setModifierConfigUrlNegativeRegexp(UrlNegativeRegexp HTTPUrlRegexp) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.UrlNegativeRegexp = UrlNegativeRegexp +} + +func (a *AppSettings) modifierConfigUrlRegexp() HTTPUrlRegexp { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.UrlRegexp +} + +func (a *AppSettings) setModifierConfigUrlRegexp(UrlRegexp HTTPUrlRegexp) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.UrlRegexp = UrlRegexp +} + +func (a *AppSettings) modifierConfigUrlRewrite() UrlRewriteMap { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.UrlRewrite +} + +func (a *AppSettings) setModifierConfigUrlRewrite(UrlRewrite UrlRewriteMap) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.UrlRewrite = UrlRewrite +} + +func (a *AppSettings) modifierConfigHeaderRewrite() HeaderRewriteMap { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.HeaderRewrite +} + +func (a *AppSettings) setModifierConfigHeaderRewrite(headerRewrite HeaderRewriteMap) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.HeaderRewrite = headerRewrite +} + +func (a *AppSettings) modifierConfigHeaderFilters() HTTPHeaderFilters { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.HeaderFilters +} + +func (a *AppSettings) setModifierConfigHeaderFilters(HeaderFilters HTTPHeaderFilters) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.HeaderFilters = HeaderFilters +} + +func (a *AppSettings) modifierConfigHeaderNegativeFilters() HTTPHeaderFilters { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.HeaderNegativeFilters +} + +func (a *AppSettings) setModifierConfigHeaderNegativeFilters(headerNegativeFilters HTTPHeaderFilters) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.HeaderNegativeFilters = headerNegativeFilters +} + +func (a *AppSettings) modifierConfigHeaderBasicAuthFilters() HTTPHeaderBasicAuthFilters { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.HeaderBasicAuthFilters +} + +func (a *AppSettings) setModifierConfigHeaderBasicAuthFilters(headerBasicAuthFilters HTTPHeaderBasicAuthFilters) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.HeaderBasicAuthFilters = headerBasicAuthFilters +} + +func (a *AppSettings) modifierConfigHeaderHashFilters() HTTPHashFilters { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.HeaderHashFilters +} + +func (a *AppSettings) setModifierConfigHeaderHashFilters(headerHashFilters HTTPHashFilters) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.HeaderHashFilters = headerHashFilters +} + +func (a *AppSettings) modifierConfigParamHashFilters() HTTPHashFilters { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.ParamHashFilters +} + +func (a *AppSettings) setModifierConfigParamHashFilters(paramHashFilters HTTPHashFilters) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.ParamHashFilters = paramHashFilters +} + +func (a *AppSettings) modifierConfigParams() HTTPParams { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.Params +} + +func (a *AppSettings) setModifierConfigParams(params HTTPParams) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.Params = params +} + +func (a *AppSettings) modifierConfigHeaders() HTTPHeaders { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.Headers +} + +func (a *AppSettings) setModifierConfigHeaders(headers HTTPHeaders) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.Headers = headers +} + +func (a *AppSettings) modifierConfigMethods() HTTPMethods { + a.RLock() + defer a.RUnlock() + return a.ModifierConfig.Methods +} + +func (a *AppSettings) setModifierConfigMethods(methods HTTPMethods) { + a.Lock() + defer a.Unlock() + a.ModifierConfig.Methods = methods +} + +func (a *AppSettings) inputKafkaConfigHost() string { + a.RLock() + defer a.RUnlock() + return a.InputKafkaConfig.Host +} + +func (a *AppSettings) setInputKafkaConfigHost(host string) { + a.Lock() + defer a.Unlock() + a.InputKafkaConfig.Host = host +} + +func (a *AppSettings) inputKafkaConfigTopic() string { + a.RLock() + defer a.RUnlock() + return a.InputKafkaConfig.Topic +} + +func (a *AppSettings) setInputKafkaConfigTopic(topic string) { + a.Lock() + defer a.Unlock() + a.InputKafkaConfig.Topic = topic +} + +func (a *AppSettings) inputKafkaConfigUseJSON() bool { + a.RLock() + defer a.RUnlock() + return a.InputKafkaConfig.UseJSON +} + +func (a *AppSettings) setInputKafkaConfigUseJSON(configJSON bool) { + a.Lock() + defer a.Unlock() + a.InputKafkaConfig.UseJSON = configJSON +} + +func (a *AppSettings) outputKafkaConfigHost() string { + a.RLock() + defer a.RUnlock() + return a.OutputKafkaConfig.Host +} + +func (a *AppSettings) setOutputKafkaConfigHost(host string) { + a.Lock() + defer a.Unlock() + a.OutputKafkaConfig.Host = host +} + +func (a *AppSettings) outputKafkaConfigTopic() string { + a.RLock() + defer a.RUnlock() + return a.OutputKafkaConfig.Topic +} + +func (a *AppSettings) setOutputKafkaConfigTopic(topic string) { + a.Lock() + defer a.Unlock() + a.OutputKafkaConfig.Topic = topic +} + +func (a *AppSettings) outputKafkaConfigUseJSON() bool { + a.RLock() + defer a.RUnlock() + return a.OutputKafkaConfig.UseJSON +} + +func (a *AppSettings) setOutputKafkaConfigUseJSON(outputKafkaConfigUseJSON bool) { + a.Lock() + defer a.Unlock() + a.OutputKafkaConfig.UseJSON = outputKafkaConfigUseJSON +} + +func (a *AppSettings) configFile() string { + a.RLock() + defer a.RUnlock() + return a.ConfigFile +} + +func (a *AppSettings) setConfigFile(configFile string) { + a.Lock() + defer a.Unlock() + a.ConfigFile = configFile +} + +func (a *AppSettings) configServerAddress() string { + a.RLock() + defer a.RUnlock() + return a.ConfigServerAddress +} + +func (a *AppSettings) setConfigServerAddress(configServerAddress string) { + a.Lock() + defer a.Unlock() + a.ConfigServerAddress = configServerAddress +} + +func (a *AppSettings) remoteConfigHost() string { + a.RLock() + defer a.RUnlock() + return a.RemoteConfigHost +} + +func (a *AppSettings) setRemoteConfigHost(remoteConfigHost string) { + a.Lock() + defer a.Unlock() + a.RemoteConfigHost = remoteConfigHost } var nestedPathMap = make(map[string]string) @@ -201,6 +1247,7 @@ func readAndUpdateConfig(body io.ReadCloser) error { if err != nil { return err } + Settings = t newConfig, err := json.Marshal(Settings) if err != nil { @@ -218,8 +1265,10 @@ func updateConfig(respBody []byte) { if err := json.Unmarshal(respBody, &t); err != nil { return } - Settings = t + Settings.Lock() + Settings = t //W newConfig, err := json.Marshal(Settings) + Settings.Unlock() if err != nil { return } @@ -239,14 +1288,18 @@ func flagz(res http.ResponseWriter, req *http.Request) { } else { viper.Set(k, v[0]) } + Settings.Lock() viper.Unmarshal(&Settings) + Settings.Unlock() } } if req.Method == "POST" { res.Write([]byte(readAndUpdateConfig(req.Body).Error())) } + Settings.RLock() data, _ := json.MarshalIndent(Settings, "", " ") + Settings.RUnlock() res.Write(data) for k, v := range req.URL.Query() { @@ -257,14 +1310,18 @@ func flagz(res http.ResponseWriter, req *http.Request) { } else { viper.Set(k, v[0]) } + Settings.Lock() viper.Unmarshal(&Settings) + Settings.Unlock() } } if req.Method == "POST" { res.Write([]byte(readAndUpdateConfig(req.Body).Error())) } + Settings.RLock() data, _ = json.MarshalIndent(Settings, "", " ") + Settings.RUnlock() res.Write(data) } @@ -449,12 +1506,14 @@ func init() { } viper.Unmarshal(&Settings) - if Settings.RemoteConfigHost != "" { + if Settings.remoteConfigHost() != "" { log.Printf("Starting to read remote config from: %s", Settings.RemoteConfigHost) go pollRemoteConfig() } - fmt.Printf("Using config: %s\n", viper.ConfigFileUsed()) + if Settings.verbose() { + fmt.Printf("Using config: %s\n", viper.ConfigFileUsed()) + } initConfigServer() @@ -469,7 +1528,8 @@ func pollRemoteConfig() { } req.Header.Set("Content-Type", "application/json") - client := &http.Client{} + client := &http.Client{Timeout: 5 * time.Second} + resp, err := client.Do(req) if err != nil { log.Printf("Error while getting config from remote server., %s", err) @@ -484,7 +1544,6 @@ func pollRemoteConfig() { } updateConfig(respBody) - resp.Body.Close() time.Sleep(time.Second) } }