Add s3 object count to trace logs (#975)
* Add s3 object count to trace logs * fix debug level
This commit is contained in:
@@ -93,7 +93,7 @@ var (
|
|||||||
syslogTLSKey = syslogScan.Flag("key", "Path to TLS key.").String()
|
syslogTLSKey = syslogScan.Flag("key", "Path to TLS key.").String()
|
||||||
syslogFormat = syslogScan.Flag("format", "Log format. Can be rfc3164 or rfc5424").String()
|
syslogFormat = syslogScan.Flag("format", "Log format. Can be rfc3164 or rfc5424").String()
|
||||||
|
|
||||||
stderrLevel = zap.NewAtomicLevel()
|
logLevel = zap.NewAtomicLevel()
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
@@ -113,9 +113,13 @@ func init() {
|
|||||||
}
|
}
|
||||||
switch {
|
switch {
|
||||||
case *trace:
|
case *trace:
|
||||||
|
log.SetLevel(5)
|
||||||
|
log.SetLevelForControl(logLevel, 5)
|
||||||
logrus.SetLevel(logrus.TraceLevel)
|
logrus.SetLevel(logrus.TraceLevel)
|
||||||
logrus.Debugf("running version %s", version.BuildVersion)
|
logrus.Debugf("running version %s", version.BuildVersion)
|
||||||
case *debug:
|
case *debug:
|
||||||
|
log.SetLevel(2)
|
||||||
|
log.SetLevelForControl(logLevel, 2)
|
||||||
logrus.SetLevel(logrus.DebugLevel)
|
logrus.SetLevel(logrus.DebugLevel)
|
||||||
logrus.Debugf("running version %s", version.BuildVersion)
|
logrus.Debugf("running version %s", version.BuildVersion)
|
||||||
default:
|
default:
|
||||||
@@ -172,11 +176,11 @@ func run(state overseer.State) {
|
|||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
logger, sync := log.New("trufflehog", log.WithConsoleSink(os.Stderr, log.WithLeveler(stderrLevel)))
|
logger, sync := log.New("trufflehog", log.WithConsoleSink(os.Stderr, log.WithLeveler(logLevel)))
|
||||||
context.SetDefaultLogger(logger)
|
ctx := context.WithLogger(context.TODO(), logger)
|
||||||
|
|
||||||
defer func() { _ = sync() }()
|
defer func() { _ = sync() }()
|
||||||
|
|
||||||
ctx := context.TODO()
|
|
||||||
e := engine.Start(ctx,
|
e := engine.Start(ctx,
|
||||||
engine.WithConcurrency(*concurrency),
|
engine.WithConcurrency(*concurrency),
|
||||||
engine.WithDecoders(decoders.DefaultDecoders()...),
|
engine.WithDecoders(decoders.DefaultDecoders()...),
|
||||||
|
|||||||
+10
-3
@@ -5,6 +5,7 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/aws/aws-sdk-go/aws"
|
"github.com/aws/aws-sdk-go/aws"
|
||||||
@@ -68,6 +69,7 @@ func (s *Source) Init(aCtx context.Context, name string, jobId, sourceId int64,
|
|||||||
s.verify = verify
|
s.verify = verify
|
||||||
s.concurrency = concurrency
|
s.concurrency = concurrency
|
||||||
s.errorCount = &sync.Map{}
|
s.errorCount = &sync.Map{}
|
||||||
|
s.log = aCtx.Logger()
|
||||||
|
|
||||||
var conn sourcespb.S3
|
var conn sourcespb.S3
|
||||||
err := anypb.UnmarshalTo(connection, &conn, proto.UnmarshalOptions{})
|
err := anypb.UnmarshalTo(connection, &conn, proto.UnmarshalOptions{})
|
||||||
@@ -137,6 +139,7 @@ func (s *Source) Chunks(ctx context.Context, chunksChan chan *sources.Chunk) err
|
|||||||
return errors.Errorf("invalid configuration given for %s source", s.name)
|
return errors.Errorf("invalid configuration given for %s source", s.name)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
objectCount := uint64(0)
|
||||||
for i, bucket := range bucketsToScan {
|
for i, bucket := range bucketsToScan {
|
||||||
if common.IsDone(ctx) {
|
if common.IsDone(ctx) {
|
||||||
return nil
|
return nil
|
||||||
@@ -165,7 +168,7 @@ func (s *Source) Chunks(ctx context.Context, chunksChan chan *sources.Chunk) err
|
|||||||
err = regionalClient.ListObjectsV2PagesWithContext(
|
err = regionalClient.ListObjectsV2PagesWithContext(
|
||||||
ctx, &s3.ListObjectsV2Input{Bucket: &bucket},
|
ctx, &s3.ListObjectsV2Input{Bucket: &bucket},
|
||||||
func(page *s3.ListObjectsV2Output, last bool) bool {
|
func(page *s3.ListObjectsV2Output, last bool) bool {
|
||||||
s.pageChunker(ctx, regionalClient, chunksChan, bucket, page, &errorCount)
|
s.pageChunker(ctx, regionalClient, chunksChan, bucket, page, &errorCount, i+1, &objectCount)
|
||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
|
|
||||||
@@ -173,13 +176,13 @@ func (s *Source) Chunks(ctx context.Context, chunksChan chan *sources.Chunk) err
|
|||||||
s.log.Error(err, "could not list objects in s3 bucket", "bucket", bucket)
|
s.log.Error(err, "could not list objects in s3 bucket", "bucket", bucket)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
s.SetProgressComplete(len(bucketsToScan), len(bucketsToScan), fmt.Sprintf("Completed scanning source %s", s.name), "")
|
s.SetProgressComplete(len(bucketsToScan), len(bucketsToScan), fmt.Sprintf("Completed scanning source %s. %d objects scanned.", s.name, objectCount), "")
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// pageChunker emits chunks onto the given channel from a page
|
// pageChunker emits chunks onto the given channel from a page
|
||||||
func (s *Source) pageChunker(ctx context.Context, client *s3.S3, chunksChan chan *sources.Chunk, bucket string, page *s3.ListObjectsV2Output, errorCount *sync.Map) {
|
func (s *Source) pageChunker(ctx context.Context, client *s3.S3, chunksChan chan *sources.Chunk, bucket string, page *s3.ListObjectsV2Output, errorCount *sync.Map, pageNumber int, objectCount *uint64) {
|
||||||
sem := semaphore.NewWeighted(int64(s.concurrency))
|
sem := semaphore.NewWeighted(int64(s.concurrency))
|
||||||
var wg sync.WaitGroup
|
var wg sync.WaitGroup
|
||||||
for _, obj := range page.Contents {
|
for _, obj := range page.Contents {
|
||||||
@@ -303,6 +306,8 @@ func (s *Source) pageChunker(ctx context.Context, client *s3.S3, chunksChan chan
|
|||||||
Verify: s.verify,
|
Verify: s.verify,
|
||||||
}
|
}
|
||||||
if handlers.HandleFile(ctx, reader, chunkSkel, chunksChan) {
|
if handlers.HandleFile(ctx, reader, chunkSkel, chunksChan) {
|
||||||
|
atomic.AddUint64(objectCount, 1)
|
||||||
|
s.log.V(5).Info("S3 object scanned.", "object_count", objectCount, "page_number", pageNumber)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -317,6 +322,8 @@ func (s *Source) pageChunker(ctx context.Context, client *s3.S3, chunksChan chan
|
|||||||
s.log.Error(err, "Could not read file data.")
|
s.log.Error(err, "Could not read file data.")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
atomic.AddUint64(objectCount, 1)
|
||||||
|
s.log.V(5).Info("S3 object scanned.", "object_count", objectCount, "page_number", pageNumber)
|
||||||
chunk.Data = chunkData
|
chunk.Data = chunkData
|
||||||
chunksChan <- &chunk
|
chunksChan <- &chunk
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user