Fix race condition while setting flags.

This commit is contained in:
arijitad
2020-08-10 16:35:35 +05:30
parent c1e61434c3
commit b47708e35d
23 changed files with 1165 additions and 106 deletions
+1 -1
View File
@@ -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()
+6 -6
View File
@@ -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")
}
+15 -15
View File
@@ -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()
+7 -7
View File
@@ -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) {
+2 -2
View File
@@ -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
}
+2 -2
View File
@@ -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() {
+2 -2
View File
@@ -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)
+1 -1
View File
@@ -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)
}
+1 -1
View File
@@ -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()
+9 -9
View File
@@ -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")
+2 -2
View File
@@ -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,
+4 -4
View File
@@ -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()
+4 -4
View File
@@ -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))
}
+6 -6
View File
@@ -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("")
}
+2 -2
View File
@@ -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()
+4 -4
View File
@@ -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)
}
+11 -11
View File
@@ -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)
+1 -1
View File
@@ -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()
}
+2 -2
View File
@@ -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]))
}
+2 -2
View File
@@ -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++ {
+13 -13
View File
@@ -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 {
+4 -4
View File
@@ -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()
+1064 -5
View File
File diff suppressed because it is too large Load Diff