From 4d1cff24785673829f1c45f94d3212ea99b947e3 Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Fri, 29 Aug 2025 16:25:09 +0530 Subject: [PATCH 1/6] fix: removing the addlsitener and removelsiterner channel --- pkg/broadcast/broadcast.go | 35 ++++++++++++++--------------------- 1 file changed, 14 insertions(+), 21 deletions(-) diff --git a/pkg/broadcast/broadcast.go b/pkg/broadcast/broadcast.go index 6f0fce6..f4698b4 100644 --- a/pkg/broadcast/broadcast.go +++ b/pkg/broadcast/broadcast.go @@ -33,31 +33,35 @@ type BroadcastServer interface { } type broadcastServer struct { - source <-chan scan.Hit - listeners []chan scan.Hit - addListener chan chan scan.Hit - removeListener chan (<-chan scan.Hit) + source <-chan scan.Hit + listeners []chan scan.Hit } // Subscribe() creates a subcribtion on broadcastServer. func (s *broadcastServer) Subscribe() <-chan scan.Hit { newListener := make(chan scan.Hit) - s.addListener <- newListener + s.listeners = append(s.listeners, newListener) + return newListener } // CancelSubscription() cancel a subcribtion on broadcastServer. func (s *broadcastServer) CancelSubscription(channel <-chan scan.Hit) { - s.removeListener <- channel + for i, ch := range s.listeners { + if ch == channel { + s.listeners[i] = s.listeners[len(s.listeners)-1] + s.listeners = s.listeners[:len(s.listeners)-1] + close(ch) + break + } + } } // NewBroadcastServer() create a broadcast server and starts new routine. func NewBroadcastServer(ctx context.Context, source <-chan scan.Hit) BroadcastServer { service := &broadcastServer{ - source: source, - listeners: make([]chan scan.Hit, 0), - addListener: make(chan chan scan.Hit, 10), - removeListener: make(chan (<-chan scan.Hit)), + source: source, + listeners: make([]chan scan.Hit, 0), } go service.serve(ctx) return service @@ -77,17 +81,6 @@ func (s *broadcastServer) serve(ctx context.Context) { select { case <-ctx.Done(): return - case newListener := <-s.addListener: - s.listeners = append(s.listeners, newListener) - case listenerToRemove := <-s.removeListener: - for i, ch := range s.listeners { - if ch == listenerToRemove { - s.listeners[i] = s.listeners[len(s.listeners)-1] - s.listeners = s.listeners[:len(s.listeners)-1] - close(ch) - break - } - } case val, ok := <-s.source: if !ok { return From 02f65afae673d58dc36964da5c8b623ee415f49e Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Fri, 29 Aug 2025 17:12:39 +0530 Subject: [PATCH 2/6] fix: creating the listeners buffered --- pkg/broadcast/broadcast.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/broadcast/broadcast.go b/pkg/broadcast/broadcast.go index f4698b4..3b1ed2d 100644 --- a/pkg/broadcast/broadcast.go +++ b/pkg/broadcast/broadcast.go @@ -39,7 +39,7 @@ type broadcastServer struct { // Subscribe() creates a subcribtion on broadcastServer. func (s *broadcastServer) Subscribe() <-chan scan.Hit { - newListener := make(chan scan.Hit) + newListener := make(chan scan.Hit, 100) s.listeners = append(s.listeners, newListener) return newListener From 7c7d0c7202f3bc53f0fa0a23b4218261b7ae0482 Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Mon, 1 Sep 2025 22:56:11 +0530 Subject: [PATCH 3/6] fix: making the listeners available first --- pkg/broadcast/broadcast.go | 43 +++++++++++++++++++++--------- pkg/core/core.go | 54 ++++++++++++++++++++------------------ pkg/writers/consoleout.go | 3 ++- pkg/writers/jsonout.go | 7 ++--- 4 files changed, 65 insertions(+), 42 deletions(-) diff --git a/pkg/broadcast/broadcast.go b/pkg/broadcast/broadcast.go index 3b1ed2d..29beeee 100644 --- a/pkg/broadcast/broadcast.go +++ b/pkg/broadcast/broadcast.go @@ -23,32 +23,35 @@ package broadcast import ( "context" + "sync" "github.com/americanexpress/earlybird/v4/pkg/scan" ) type BroadcastServer interface { - Subscribe() <-chan scan.Hit CancelSubscription(<-chan scan.Hit) + GetListeners() []chan scan.Hit } type broadcastServer struct { source <-chan scan.Hit listeners []chan scan.Hit + wg *sync.WaitGroup } -// Subscribe() creates a subcribtion on broadcastServer. +// Subscribe() creates a subscription on broadcastServer. func (s *broadcastServer) Subscribe() <-chan scan.Hit { - newListener := make(chan scan.Hit, 100) + newListener := make(chan scan.Hit) s.listeners = append(s.listeners, newListener) return newListener } -// CancelSubscription() cancel a subcribtion on broadcastServer. +// CancelSubscription() cancel a subscription on broadcastServer. func (s *broadcastServer) CancelSubscription(channel <-chan scan.Hit) { for i, ch := range s.listeners { if ch == channel { + s.wg.Done() s.listeners[i] = s.listeners[len(s.listeners)-1] s.listeners = s.listeners[:len(s.listeners)-1] close(ch) @@ -57,25 +60,41 @@ func (s *broadcastServer) CancelSubscription(channel <-chan scan.Hit) { } } +func (s *broadcastServer) GetListeners() []chan scan.Hit { + return s.listeners +} + +func (s *broadcastServer) AddSubscriber(count int, wg *sync.WaitGroup) { + for i := 0; i < count; i++ { + s.wg.Add(1) + s.Subscribe() + } +} + +func (s *broadcastServer) CloseBroadcast() { + for _, listener := range s.listeners { + if listener != nil { + close(listener) + } + } + // defer s.wg.Done() +} + // NewBroadcastServer() create a broadcast server and starts new routine. -func NewBroadcastServer(ctx context.Context, source <-chan scan.Hit) BroadcastServer { +func NewBroadcastServer(ctx context.Context, source <-chan scan.Hit, count int, wg *sync.WaitGroup) BroadcastServer { service := &broadcastServer{ source: source, listeners: make([]chan scan.Hit, 0), + wg: wg, } + service.AddSubscriber(count, wg) go service.serve(ctx) return service } // serve() run the server and manages listener counts. func (s *broadcastServer) serve(ctx context.Context) { - defer func() { - for _, listener := range s.listeners { - if listener != nil { - close(listener) - } - } - }() + defer s.CloseBroadcast() for { select { diff --git a/pkg/core/core.go b/pkg/core/core.go index ce7934a..1c5caaf 100644 --- a/pkg/core/core.go +++ b/pkg/core/core.go @@ -293,11 +293,16 @@ func (eb *EarlybirdCfg) Scan() { if err != nil { log.Fatal("Failed to get FileContext: ", err) } + var wg sync.WaitGroup HitChannel := make(chan scan.Hit) + ctx, cancel := context.WithCancel(context.Background()) + + // Send output to a writer | creating the hit results receiver first. + eb.WriteResults(start, HitChannel, fileContext, ctx, &wg) go scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) - // Send output to a writer - eb.WriteResults(start, HitChannel, fileContext) + wg.Wait() + cancel() utils.DeleteGit(eb.Config.Gitrepo, eb.Config.SearchDir) if eb.Config.FailScan { @@ -336,40 +341,37 @@ func (eb *EarlybirdCfg) FileContext() (fileContext file.Context, err error) { } // WriteResults reads hits from the channel to the console or target file -func (eb *EarlybirdCfg) WriteResults(start time.Time, HitChannel chan scan.Hit, fileContext file.Context) { +func (eb *EarlybirdCfg) WriteResults(start time.Time, HitChannel chan scan.Hit, fileContext file.Context, ctx context.Context, wg *sync.WaitGroup) { // Send output to a writer var err error - - if eb.Config.WithConsole && eb.Config.OutputFormat == "json" { - var wg sync.WaitGroup - wg.Add(2) - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - broadcaster := broadcast.NewBroadcastServer(ctx, HitChannel) - listener1 := broadcaster.Subscribe() - listener2 := broadcaster.Subscribe() + // + if true { + // initializing the broadcaster with two listeners, and starting the broadcast server + broadcaster := broadcast.NewBroadcastServer(ctx, HitChannel, 2, wg) + listener := broadcaster.GetListeners() go func() { defer wg.Done() - err = writers.WriteConsole(listener1, "", eb.Config.ShowFullLine) + err = writers.WriteConsole(listener[0], "", eb.Config.ShowFullLine) log.Printf("\n%d files scanned in %s", len(fileContext.Files), time.Since(start)) - log.Printf("\n%d rules observed\n", len(scan.CombinedRules)) }() go func() { defer wg.Done() - err = writers.WriteJSON(listener2, eb.Config, fileContext, eb.Config.OutputFile) + err = writers.WriteJSON(listener[1], eb.Config, fileContext, eb.Config.OutputFile) }() - wg.Wait() } else { - switch { - case eb.Config.OutputFormat == "json": - err = writers.WriteJSON(HitChannel, eb.Config, fileContext, eb.Config.OutputFile) - case eb.Config.OutputFormat == "csv": - err = writers.WriteCSV(HitChannel, eb.Config.OutputFile) - default: - err = writers.WriteConsole(HitChannel, eb.Config.OutputFile, eb.Config.ShowFullLine) - log.Printf("\n%d files scanned in %s", len(fileContext.Files), time.Since(start)) - log.Printf("\n%d rules observed\n", len(scan.CombinedRules)) - } + wg.Add(1) + go func() { + defer wg.Done() + switch { + case eb.Config.OutputFormat == "json": + err = writers.WriteJSON(HitChannel, eb.Config, fileContext, eb.Config.OutputFile) + case eb.Config.OutputFormat == "csv": + err = writers.WriteCSV(HitChannel, eb.Config.OutputFile) + default: + err = writers.WriteConsole(HitChannel, eb.Config.OutputFile, eb.Config.ShowFullLine) + log.Printf("\n%d files scanned in %s", len(fileContext.Files), time.Since(start)) + } + }() } if err != nil { log.Println("Writing Results failed:", err) diff --git a/pkg/writers/consoleout.go b/pkg/writers/consoleout.go index 3ad9273..6f5b96b 100644 --- a/pkg/writers/consoleout.go +++ b/pkg/writers/consoleout.go @@ -35,7 +35,7 @@ type issue struct { var issues = make(map[string]int) -//WriteConsole streams hits from the result channel to the command line or target file +// WriteConsole streams hits from the result channel to the command line or target file func WriteConsole(hits <-chan scan.Hit, fileName string, showFullLine bool) error { // If no filename was passed in, just print to stdout if fileName == "" { @@ -53,6 +53,7 @@ func WriteConsole(hits <-chan scan.Hit, fileName string, showFullLine bool) erro } } displayIssues() + log.Printf("\n%d rules observed\n", len(scan.CombinedRules)) return nil } diff --git a/pkg/writers/jsonout.go b/pkg/writers/jsonout.go index d1e89f8..8be81b5 100755 --- a/pkg/writers/jsonout.go +++ b/pkg/writers/jsonout.go @@ -19,11 +19,12 @@ package writers import ( "encoding/json" "fmt" + "os" + "time" + cfgReader "github.com/americanexpress/earlybird/v4/pkg/config" "github.com/americanexpress/earlybird/v4/pkg/file" "github.com/americanexpress/earlybird/v4/pkg/scan" - "os" - "time" ) // WriteJSON takes the hits, converts them into JSON report and passing report to reportToJSONWriter(). @@ -52,7 +53,7 @@ func WriteJSON(hits <-chan scan.Hit, config cfgReader.EarlybirdConfig, fileConte return err } -//reportToJSONWriter Outputs an object as a JSON blob to an output file or console +// reportToJSONWriter Outputs an object as a JSON blob to an output file or console func reportToJSONWriter(v interface{}, fileName string) (s string, err error) { b, err := json.MarshalIndent(v, "", "\t") if err != nil { From ddc5e11816610330118eb79a29533e93eac4dbca Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Mon, 1 Sep 2025 23:59:16 +0530 Subject: [PATCH 4/6] fix: removing context from server --- pkg/broadcast/broadcast.go | 34 +++++++++------------------------- pkg/core/core.go | 19 +++++++++++-------- 2 files changed, 20 insertions(+), 33 deletions(-) diff --git a/pkg/broadcast/broadcast.go b/pkg/broadcast/broadcast.go index 29beeee..c1573c7 100644 --- a/pkg/broadcast/broadcast.go +++ b/pkg/broadcast/broadcast.go @@ -22,7 +22,6 @@ package broadcast import ( - "context" "sync" "github.com/americanexpress/earlybird/v4/pkg/scan" @@ -77,43 +76,28 @@ func (s *broadcastServer) CloseBroadcast() { close(listener) } } - // defer s.wg.Done() } // NewBroadcastServer() create a broadcast server and starts new routine. -func NewBroadcastServer(ctx context.Context, source <-chan scan.Hit, count int, wg *sync.WaitGroup) BroadcastServer { +func NewBroadcastServer(source <-chan scan.Hit, count int, wg *sync.WaitGroup) BroadcastServer { service := &broadcastServer{ source: source, listeners: make([]chan scan.Hit, 0), wg: wg, } service.AddSubscriber(count, wg) - go service.serve(ctx) + go service.broadCastData() return service } -// serve() run the server and manages listener counts. -func (s *broadcastServer) serve(ctx context.Context) { - defer s.CloseBroadcast() - - for { - select { - case <-ctx.Done(): - return - case val, ok := <-s.source: - if !ok { - return - } - for _, listener := range s.listeners { - if listener != nil { - select { - case listener <- val: - case <-ctx.Done(): - return - } - - } +// broadCastData() run the server and manages listener counts. +func (s *broadcastServer) broadCastData() { + for val := range s.source { + for _, listener := range s.listeners { + if listener != nil { + listener <- val } } } + s.CloseBroadcast() } diff --git a/pkg/core/core.go b/pkg/core/core.go index 1c5caaf..7cf0271 100644 --- a/pkg/core/core.go +++ b/pkg/core/core.go @@ -18,7 +18,6 @@ package core import ( "bufio" - "context" "flag" "fmt" "log" @@ -295,14 +294,12 @@ func (eb *EarlybirdCfg) Scan() { } var wg sync.WaitGroup HitChannel := make(chan scan.Hit) - ctx, cancel := context.WithCancel(context.Background()) // Send output to a writer | creating the hit results receiver first. - eb.WriteResults(start, HitChannel, fileContext, ctx, &wg) - go scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) + eb.WriteResults(start, HitChannel, fileContext, &wg) + scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) wg.Wait() - cancel() utils.DeleteGit(eb.Config.Gitrepo, eb.Config.SearchDir) if eb.Config.FailScan { @@ -341,22 +338,24 @@ func (eb *EarlybirdCfg) FileContext() (fileContext file.Context, err error) { } // WriteResults reads hits from the channel to the console or target file -func (eb *EarlybirdCfg) WriteResults(start time.Time, HitChannel chan scan.Hit, fileContext file.Context, ctx context.Context, wg *sync.WaitGroup) { +func (eb *EarlybirdCfg) WriteResults(start time.Time, HitChannel chan scan.Hit, fileContext file.Context, wg *sync.WaitGroup) { // Send output to a writer var err error // - if true { + if eb.Config.WithConsole && eb.Config.OutputFormat == "json" { // initializing the broadcaster with two listeners, and starting the broadcast server - broadcaster := broadcast.NewBroadcastServer(ctx, HitChannel, 2, wg) + broadcaster := broadcast.NewBroadcastServer(HitChannel, 2, wg) listener := broadcaster.GetListeners() go func() { defer wg.Done() err = writers.WriteConsole(listener[0], "", eb.Config.ShowFullLine) log.Printf("\n%d files scanned in %s", len(fileContext.Files), time.Since(start)) + printError(err) }() go func() { defer wg.Done() err = writers.WriteJSON(listener[1], eb.Config, fileContext, eb.Config.OutputFile) + printError(err) }() } else { wg.Add(1) @@ -371,8 +370,12 @@ func (eb *EarlybirdCfg) WriteResults(start time.Time, HitChannel chan scan.Hit, err = writers.WriteConsole(HitChannel, eb.Config.OutputFile, eb.Config.ShowFullLine) log.Printf("\n%d files scanned in %s", len(fileContext.Files), time.Since(start)) } + printError(err) }() } +} + +func printError(err error) { if err != nil { log.Println("Writing Results failed:", err) } From 1485ed2a84553966be792b41a1770d672ef9c5a8 Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Tue, 2 Sep 2025 09:22:13 +0530 Subject: [PATCH 5/6] chore: adding the comment --- pkg/core/core.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/core/core.go b/pkg/core/core.go index 7cf0271..30120e1 100644 --- a/pkg/core/core.go +++ b/pkg/core/core.go @@ -299,7 +299,7 @@ func (eb *EarlybirdCfg) Scan() { eb.WriteResults(start, HitChannel, fileContext, &wg) scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) - wg.Wait() + wg.Wait() // this wait ensures that all writers goroutine are finished utils.DeleteGit(eb.Config.Gitrepo, eb.Config.SearchDir) if eb.Config.FailScan { From 7e2e42186493f5360df9f0c2681c776c08ab0f64 Mon Sep 17 00:00:00 2001 From: Afham Fardeen Date: Tue, 2 Sep 2025 10:46:36 +0530 Subject: [PATCH 6/6] chore: adding commits --- pkg/core/core.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/pkg/core/core.go b/pkg/core/core.go index 30120e1..5ebbbac 100644 --- a/pkg/core/core.go +++ b/pkg/core/core.go @@ -295,9 +295,8 @@ func (eb *EarlybirdCfg) Scan() { var wg sync.WaitGroup HitChannel := make(chan scan.Hit) - // Send output to a writer | creating the hit results receiver first. - eb.WriteResults(start, HitChannel, fileContext, &wg) - scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) + eb.WriteResults(start, HitChannel, fileContext, &wg) // Registering the hit receiver. + scan.SearchFiles(&eb.Config, fileContext.Files, fileContext.CompressPaths, fileContext.ConvertPaths, HitChannel) // sending the hits to the channel from the worker threads running on go-routine. wg.Wait() // this wait ensures that all writers goroutine are finished