Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
2 changes: 2 additions & 0 deletions services/gtc/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ module github.com/emoss08/gtc
go 1.26

require (
github.com/alicebob/miniredis/v2 v2.38.0
github.com/bytedance/sonic v1.15.2
github.com/go-chi/chi/v5 v5.2.5
github.com/go-ozzo/ozzo-validation/v4 v4.4.1
Expand Down Expand Up @@ -40,6 +41,7 @@ require (
github.com/prometheus/procfs v0.21.1 // indirect
github.com/rogpeppe/go-internal v1.14.1 // indirect
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
github.com/yuin/gopher-lua v1.1.1 // indirect
go.uber.org/atomic v1.11.0 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.yaml.in/yaml/v3 v3.0.5 // indirect
Expand Down
26 changes: 26 additions & 0 deletions services/gtc/go.sum
Original file line number Diff line number Diff line change
@@ -1,4 +1,8 @@
github.com/alicebob/miniredis/v2 v2.38.0 h1:nZAzCR+Lj+Vxk4ZXzm2NuKq2O33RXj1XxJ2e2uP9jiw=
github.com/alicebob/miniredis/v2 v2.38.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM=
github.com/andybalholm/brotli v1.2.2 h1:HzTuoo2ErYQqf5qvcJInB8uvqSVxRttzkFexPWtnceM=
github.com/andybalholm/brotli v1.2.2/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
github.com/asaskevich/govalidator v0.0.0-20210307081110-f21760c49a8d/go.mod h1:WaHUgvxTVq04UNunO+XhnAqY/wQc+bxr74GqbsZ/Jqw=
github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2 h1:DklsrG3dyBCFEj5IhUbnKptjxatkF07cF2ak3yi77so=
github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2/go.mod h1:WaHUgvxTVq04UNunO+XhnAqY/wQc+bxr74GqbsZ/Jqw=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
Expand All @@ -10,18 +14,22 @@ github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0
github.com/bytedance/gopkg v0.1.4 h1:oZnQwnX82KAIWb7033bEwtxvTqXcYMxDBaQxo5JJHWM=
github.com/bytedance/gopkg v0.1.4/go.mod h1:v1zWfPm21Fb+OsyXN2VAHdL6TBb2L88anLQgdyje6R4=
github.com/bytedance/sonic v1.15.2 h1:90H+rcF/FwLXwfB1cudOLq/je83n683Utf4Cbp0xHCo=
github.com/bytedance/sonic v1.15.2/go.mod h1:mT2NbXunuaEbnZ+mRIX/vYqKISmgEuHFDI4UzmKx2SA=
github.com/bytedance/sonic/loader v0.5.2 h1:0QtP1gevc1OZ6/H8Lb9BRZiCXd1Ftjd3OKuj1T1lBIo=
github.com/bytedance/sonic/loader v0.5.2/go.mod h1:AR4NYCk5DdzZizZ5djGqQ92eEhCCcdf5x77udYiSJRo=
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/cloudwego/base64x v0.1.7 h1:NppS+Fgzg5ovhn4NkUXaDT3x9jldgH5ToMCqzBSi2zI=
github.com/cloudwego/base64x v0.1.7/go.mod h1:Cu1PV9zfrSf7ET2tIbWbbEy7jO7HHJ13q4X2SQ8aWYg=
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug=
github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0=
github.com/go-ozzo/ozzo-validation/v4 v4.4.1 h1:AQ3X8zHnXEuNE04pyc1H/nmIlroNjgZ7hcY7Xv/IgH8=
github.com/go-ozzo/ozzo-validation/v4 v4.4.1/go.mod h1:4ZtPNefSnNq39wjL+2We8y2ysqEX/S4D5mPybufHd7Y=
github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY=
github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
Expand All @@ -35,30 +43,40 @@ github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5ey
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.10.0 h1:VhSvgU2jSli8o3AqIEOTJr7rZwAEUVo4E4XhR94Zfr0=
github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
github.com/klauspost/compress v1.19.2 h1:hMRETovs/pu/dVWN7zIT1PGG8t509MwT6bO7XSi26R8=
github.com/klauspost/compress v1.19.2/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/cpuid/v2 v2.4.0 h1:S6Hrbc7+ywsr0r+RLapfGBHfyefhCTwEh3A0tV913Dw=
github.com/klauspost/cpuid/v2 v2.4.0/go.mod h1:19jmZ9mjzoF//ddRSUsv0zfBTJWh3QJh9FNxZTMrGxU=
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
github.com/meilisearch/meilisearch-go v0.36.3 h1:Yx1aTY5jDgtbStPVkhJTDoLnZTy5sejQSPyjfNMy6e4=
github.com/meilisearch/meilisearch-go v0.36.3/go.mod h1:hWcR0MuWLSzHfbz9GGzIr3s9rnXLm1jqkmHkJPbUSvM=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA=
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU=
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
github.com/redis/go-redis/v9 v9.22.0 h1:laDvpYXTJtZLloinw1fA5Kqd6HAEH2XKxOkG/PDq2F0=
github.com/redis/go-redis/v9 v9.22.0/go.mod h1:y2g0Wj8rQvuK0ELM+oxSudcLtC09JScs98I/X9gRWY4=
github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs=
github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ=
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
Expand All @@ -78,7 +96,10 @@ github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS
github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08=
github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU=
github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E=
github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M=
github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw=
github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs=
github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s=
go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE=
go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
Expand All @@ -90,10 +111,15 @@ go.uber.org/zap v1.28.0/go.mod h1:rDLpOi171uODNm/mxFcuYWxDsqWSAVkFdX4XojSKg/Q=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
golang.org/x/arch v0.29.0 h1:8sSET5wB0+exBm0FGmOtdHMqjlRdV2DRD3/IV6OZgho=
golang.org/x/arch v0.29.0/go.mod h1:0X+GdSIP+kL5wPmpK7sdkEVTt2XoYP0cSjQSbZBwOi8=
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs=
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY=
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
Expand Down
38 changes: 33 additions & 5 deletions services/gtc/internal/adapters/secondary/meilisearch/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,18 @@ func (s *Sink) Write(ctx context.Context, projection domain.Projection, record d
return err
}

if record.Operation == domain.OperationTruncate {
task, err := index.DeleteDocumentsByFilterWithContext(
ctx,
truncateFilter(projection),
nil,
)
if err != nil {
return fmt.Errorf("truncate projection %s: %w", projection.Name, err)
}
return s.waitForTask(ctx, projection, record, task, "truncate projection documents")
}

keyField, key, err := documentKey(record, projection.PrimaryKeys)
if err != nil {
return err
Expand Down Expand Up @@ -139,11 +151,7 @@ func (s *Sink) ensureFilterableAttributes(
index meilisearch.IndexManager,
projection domain.Projection,
) error {
if len(projection.FilterableFields) == 0 {
return nil
}

filterableFields := append([]string(nil), projection.FilterableFields...)
filterableFields := filterableFieldsFor(projection)
filterableSettings := make([]any, 0, len(filterableFields))
for _, field := range filterableFields {
filterableSettings = append(filterableSettings, field)
Expand Down Expand Up @@ -179,6 +187,26 @@ func (s *Sink) ensureFilterableAttributes(
return nil
}

func filterableFieldsFor(projection domain.Projection) []string {
fields := make([]string, 0, len(projection.FilterableFields)+2)
fields = append(fields, projection.FilterableFields...)
for _, meta := range []string{"_projection", "_source_table"} {
if !slices.Contains(fields, meta) {
fields = append(fields, meta)
}
}

return fields
}

func truncateFilter(projection domain.Projection) string {
return fmt.Sprintf(
"_projection = %q AND _source_table = %q",
projection.Name,
projection.FullTableName(),
)
}

func primaryKey(record domain.SourceRecord, keyFields []string) (string, error) {
values, err := domain.PrimaryKey(record, keyFields)
if err != nil {
Expand Down
53 changes: 53 additions & 0 deletions services/gtc/internal/adapters/secondary/meilisearch/sink_test.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,64 @@
package meilisearch

import (
"slices"
"testing"

"github.com/emoss08/gtc/internal/core/domain"
)

func TestFilterableFieldsForAppendsMetadataFields(t *testing.T) {
t.Parallel()

fields := filterableFieldsFor(domain.Projection{
FilterableFields: []string{"organization_id", "business_unit_id"},
})

want := []string{"organization_id", "business_unit_id", "_projection", "_source_table"}
if !slices.Equal(fields, want) {
t.Fatalf("expected filterable fields %v, got %v", want, fields)
}
}

func TestFilterableFieldsForCoversProjectionsWithoutConfiguredFields(t *testing.T) {
t.Parallel()

fields := filterableFieldsFor(domain.Projection{})

want := []string{"_projection", "_source_table"}
if !slices.Equal(fields, want) {
t.Fatalf("expected filterable fields %v, got %v", want, fields)
}
}

func TestFilterableFieldsForDoesNotDuplicateMetadataFields(t *testing.T) {
t.Parallel()

fields := filterableFieldsFor(domain.Projection{
FilterableFields: []string{"_projection", "status"},
})

want := []string{"_projection", "status", "_source_table"}
if !slices.Equal(fields, want) {
t.Fatalf("expected filterable fields %v, got %v", want, fields)
}
}

func TestTruncateFilterScopesToProjectionAndSourceTable(t *testing.T) {
t.Parallel()

filter := truncateFilter(domain.Projection{
Name: "shipment-search",
SourceSchema: "public",
SourceTable: "shipments",
})

want := `_projection = "shipment-search" AND _source_table = "public.shipments"`
if filter != want {
t.Fatalf("expected filter %s, got %s", want, filter)
}
}

func TestDocumentKeyPrefersIDField(t *testing.T) {
t.Parallel()

Expand Down
110 changes: 100 additions & 10 deletions services/gtc/internal/adapters/secondary/redis/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,11 @@ import (
"go.uber.org/zap"
)

const (
truncateScanCount = 512
truncateDeleteBatch = 256
)

type baseSink struct {
client *goredis.Client
logger *zap.Logger
Expand Down Expand Up @@ -86,6 +91,10 @@ func (s *JSONSink) Shutdown(ctx context.Context) error {
}

func (s *JSONSink) writeJSON(ctx context.Context, projection domain.Projection, record domain.SourceRecord) error {
if record.Operation == domain.OperationTruncate {
return s.truncateJSON(ctx, projection, record)
}

key, err := s.renderTemplate(projection.Name, projection.Destination.KeyTemplate, projection.PrimaryKeys, record)
if err != nil {
return err
Expand Down Expand Up @@ -127,6 +136,77 @@ func (s *JSONSink) writeJSON(ctx context.Context, projection domain.Projection,
return nil
}

func (s *JSONSink) truncateJSON(
ctx context.Context,
projection domain.Projection,
record domain.SourceRecord,
) error {
tmpl, err := s.template(projection.Name, projection.Destination.KeyTemplate)
if err != nil {
return err
}

pattern, err := tmpl.WildcardPattern(record, projection.PrimaryKeys)
if err != nil {
return fmt.Errorf("truncate projection %s: %w", projection.Name, err)
}

deleted, err := s.deleteMatchingKeys(ctx, pattern)
if err != nil {
return fmt.Errorf(
"truncate projection %s: delete keys matching %q: %w",
projection.Name,
pattern,
err,
)
}

s.logger.Info("truncated redis json projection",
zap.String("projection", projection.Name),
zap.String("table", record.FullTableName()),
zap.String("pattern", pattern),
zap.Int64("deleted_keys", deleted),
)

return nil
}

func (s *baseSink) deleteMatchingKeys(ctx context.Context, pattern string) (int64, error) {
var deleted int64
batch := make([]string, 0, truncateDeleteBatch)

flush := func() error {
if len(batch) == 0 {
return nil
}
removed, err := s.client.Unlink(ctx, batch...).Result()
if err != nil {
return err
}
deleted += removed
batch = batch[:0]
return nil
}

iter := s.client.Scan(ctx, 0, pattern, truncateScanCount).Iterator()
for iter.Next(ctx) {
batch = append(batch, iter.Val())
if len(batch) >= truncateDeleteBatch {
if err := flush(); err != nil {
return deleted, err
}
}
}
if err := iter.Err(); err != nil {
return deleted, err
}
if err := flush(); err != nil {
return deleted, err
}

return deleted, nil
}

func (s *StreamSink) Kind() domain.DestinationKind {
return domain.DestinationRedisStream
}
Expand Down Expand Up @@ -174,25 +254,35 @@ func (s *StreamSink) writeStream(ctx context.Context, projection domain.Projecti
}

func (s *baseSink) renderTemplate(name string, pattern string, primaryKeys []string, record domain.SourceRecord) (string, error) {
tmpl, err := s.template(name, pattern)
if err != nil {
return "", err
}

return tmpl.Execute(record, primaryKeys)
}

func (s *baseSink) template(name string, pattern string) (*Template, error) {
key := name + "::" + pattern

s.mu.RLock()
tmpl, ok := s.templates[key]
s.mu.RUnlock()

if !ok {
parsed, err := ParseTemplate(pattern)
if err != nil {
return "", err
}
if ok {
return tmpl, nil
}

s.mu.Lock()
s.templates[key] = parsed
s.mu.Unlock()
tmpl = parsed
parsed, err := ParseTemplate(pattern)
if err != nil {
return nil, err
}

return tmpl.Execute(record, primaryKeys)
s.mu.Lock()
s.templates[key] = parsed
s.mu.Unlock()

return parsed, nil
}

func streamArgs(stream string, payload string) *goredis.XAddArgs {
Expand Down
Loading
Loading