Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
"github.com/samber/lo"

"github.com/rudderlabs/rudder-go-kit/config"
"github.com/rudderlabs/rudder-go-kit/filemanager"
"github.com/rudderlabs/rudder-go-kit/logger"
"github.com/rudderlabs/rudder-go-kit/stats"

Expand Down Expand Up @@ -80,7 +81,13 @@

handle.loggedEvents = 0
handle.loggedEventsMu = sync.Mutex{}
handle.loggedFileName = generateLogFileName()

var err error
handle.samplingFileManager, err = getSamplingUploader(conf, log)
if err != nil {
log.Errorn("failed to create dt sampling file manager", obskit.Error(err))
handle.samplingFileManager = nil
}

Check warning on line 90 in processor/internal/transformer/destination_transformer/destination_transformer.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/destination_transformer.go#L88-L90

Added lines #L88 - L90 were not covered by tests

handle.config.compactionEnabled = conf.GetReloadableBoolVar(false, "Processor.DestinationTransformer.compactionEnabled", "Transformer.compactionEnabled")

Expand Down Expand Up @@ -116,9 +123,9 @@
mismatchedEvents stats.Counter
}

loggedEventsMu sync.Mutex
loggedEvents int64
loggedFileName string
loggedEventsMu sync.Mutex
loggedEvents int64
samplingFileManager *filemanager.S3Manager
}

func (d *Client) transform(ctx context.Context, clientEvents []types.TransformerEvent) types.Response {
Expand Down Expand Up @@ -391,7 +398,9 @@
legacyTransformerResponse := c.transform(ctx, clientEvents)
embeddedTransformerResponse := impl(ctx, clientEvents)

c.CompareAndLog(embeddedTransformerResponse, legacyTransformerResponse)
go func() {
c.CompareAndLog(ctx, embeddedTransformerResponse, legacyTransformerResponse)
}()

Check warning on line 403 in processor/internal/transformer/destination_transformer/destination_transformer.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/destination_transformer.go#L401-L403

Added lines #L401 - L403 were not covered by tests

return legacyTransformerResponse
}
Expand Down Expand Up @@ -429,3 +438,32 @@
}
return jsonrs.Marshal(data)
}

func getSamplingUploader(conf *config.Config, log logger.Logger) (*filemanager.S3Manager, error) {
var (
bucket = conf.GetString("DTSampling.Bucket", "processor-dt-sampling")
endpoint = conf.GetString("DTSampling.Endpoint", "")
accessKeyID = conf.GetStringVar("", "DTSampling.AccessKeyId", "AWS_ACCESS_KEY_ID")
accessKey = conf.GetStringVar("", "DTSampling.AccessKey", "AWS_SECRET_ACCESS_KEY")
s3ForcePathStyle = conf.GetBool("DTSampling.S3ForcePathStyle", false)
disableSSL = conf.GetBool("DTSampling.DisableSsl", false)
enableSSE = conf.GetBoolVar(false, "DTSampling.EnableSse", "AWS_ENABLE_SSE")
useGlue = conf.GetBool("DTSampling.UseGlue", false)
region = conf.GetStringVar("us-east-1", "DTSampling.Region", "AWS_DEFAULT_REGION")
)
s3Config := map[string]any{
"bucketName": bucket,
"endpoint": endpoint,
"accessKeyID": accessKeyID,
"accessKey": accessKey,
"s3ForcePathStyle": s3ForcePathStyle,
"disableSSL": disableSSL,
"enableSSE": enableSSE,
"useGlue": useGlue,
"region": region,
}

return filemanager.NewS3Manager(s3Config, log.Withn(logger.NewStringField("component", "dt-uploader")), func() time.Duration {
return conf.GetDuration("DTSampling.Timeout", 120, time.Second)
})

Check warning on line 468 in processor/internal/transformer/destination_transformer/destination_transformer.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/destination_transformer.go#L467-L468

Added lines #L467 - L468 were not covered by tests
}
59 changes: 28 additions & 31 deletions processor/internal/transformer/destination_transformer/logger.go
Original file line number Diff line number Diff line change
@@ -1,23 +1,31 @@
package destination_transformer

import (
"bytes"
"context"
"fmt"
"path"

"github.com/google/go-cmp/cmp"
"github.com/google/uuid"
"github.com/samber/lo"

"github.com/rudderlabs/rudder-go-kit/config"
"github.com/rudderlabs/rudder-go-kit/logger"
obskit "github.com/rudderlabs/rudder-observability-kit/go/labels"

"github.com/rudderlabs/rudder-go-kit/stringify"

"github.com/rudderlabs/rudder-server/jsonrs"
"github.com/rudderlabs/rudder-server/processor/types"
"github.com/rudderlabs/rudder-server/utils/misc"
)

func (c *Client) CompareAndLog(
ctx context.Context,
embeddedResponse, legacyResponse types.Response,
) {
if c.samplingFileManager == nil { // Cannot upload, we should just report the issue with no diff
c.log.Errorn("DestinationTransformer sanity check failed")
return
}

Check warning on line 27 in processor/internal/transformer/destination_transformer/logger.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/logger.go#L24-L27

Added lines #L24 - L27 were not covered by tests

c.loggedEventsMu.Lock()
defer c.loggedEventsMu.Unlock()

Expand All @@ -29,20 +37,28 @@

differingResponse, sampleDiff := c.differingEvents(embeddedResponse, legacyResponse)
if len(differingResponse) == 0 && sampleDiff == "" {
c.log.Infof("Embedded and legacy responses are matches")
return
}

logEntries := lo.Map(differingResponse, func(item types.TransformerResponse, index int) string {
return stringify.Any(item)
})
if err := c.write(append([]string{sampleDiff}, logEntries...)); err != nil {
c.log.Warnn("Error logging events", obskit.Error(err))
objName := path.Join("embedded-dt-samples", config.GetKubeNamespace(), uuid.New().String())
differingResponseJSON, err := jsonrs.Marshal(differingResponse)
if err != nil {
c.log.Errorn("DestinationTransformer sanity check failed (cannot encode differingResponse)", obskit.Error(err))
return
}

Check warning on line 48 in processor/internal/transformer/destination_transformer/logger.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/logger.go#L43-L48

Added lines #L43 - L48 were not covered by tests

// upload sample diff and differing response to s3
file, err := c.samplingFileManager.UploadReader(ctx, objName, bytes.NewReader(append([]byte(sampleDiff), differingResponseJSON...)))
if err != nil {
c.log.Errorn("Error uploading DestinationTransformer sanity check diff file", obskit.Error(err))

Check warning on line 53 in processor/internal/transformer/destination_transformer/logger.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/logger.go#L51-L53

Added lines #L51 - L53 were not covered by tests
return
}

c.log.Infof("Successfully logged events: %d", len(logEntries))
c.loggedEvents += int64(len(logEntries))
c.log.Errorn("DestinationTransformer sanity check failed",
logger.NewStringField("location", file.Location),
logger.NewStringField("objectName", file.ObjectName),
)
c.loggedEvents += int64(len(differingResponse))

Check warning on line 61 in processor/internal/transformer/destination_transformer/logger.go

View check run for this annotation

Codecov / codecov/patch

processor/internal/transformer/destination_transformer/logger.go#L57-L61

Added lines #L57 - L61 were not covered by tests
}

func (c *Client) differingEvents(
Expand Down Expand Up @@ -95,22 +111,3 @@
c.stats.mismatchedEvents.Count(differedEventsCount)
return differedSampleEvents, sampleDiff
}

func (c *Client) write(data []string) error {
writer, err := misc.CreateGZ(c.loggedFileName)
if err != nil {
return fmt.Errorf("creating buffered writer: %w", err)
}
defer func() { _ = writer.Close() }()

for _, entry := range data {
if _, err := writer.Write([]byte(entry + "\n")); err != nil {
return fmt.Errorf("writing log entry: %w", err)
}
}
return nil
}

func generateLogFileName() string {
return fmt.Sprintf("destination_transformations_debug_%s.log.gz", uuid.NewString())
}
Loading