Skip to content
Merged
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
14 changes: 14 additions & 0 deletions pipeline/metadata/hash.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package metadata

import (
"github.com/cespare/xxhash/v2"
)

func Hash(parts ...string) uint64 {
h := xxhash.New()
for _, p := range parts {
_, _ = h.WriteString(p)
_, _ = h.WriteString("|")
}
return h.Sum64()
}
62 changes: 62 additions & 0 deletions pipeline/metadata/hash_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
package metadata

import (
"testing"

"github.com/stretchr/testify/assert"
)

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

tests := []struct {
name string
parts []string
}{
{
name: "single part",
parts: []string{"hello"},
},
{
name: "multiple parts",
parts: []string{"hello", "world", "foo"},
},
{
name: "empty",
parts: []string{},
},
{
name: "empty parts",
parts: []string{"", "a"},
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()

result := Hash(tt.parts...)
assert.NotEmpty(t, result)
})
}
}

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

assert.Equal(t,
Hash("a", "b"),
Hash("a", "b"),
"same inputs should produce the same hash",
)
}

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

assert.NotEqual(t,
Hash("a"),
Hash("b"),
"different inputs should produce different hashes",
)
}
41 changes: 4 additions & 37 deletions pipeline/metadata/templater.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import (
"bytes"
"fmt"
"regexp"
"strconv"
"strings"
"sync"
"text/template"
Expand Down Expand Up @@ -53,7 +52,7 @@ type MetaTemplater struct {
valueTypes *orderedmap.OrderedMap[string, ValueType]
poolBuffer sync.Pool
logger *zap.Logger
cache *lru.Cache[string, MetaData]
cache *lru.Cache[uint64, MetaData]
}

func NewMetaTemplater(templates cfg.MetaTemplates, logger *zap.Logger, cacheSize int) *MetaTemplater {
Expand Down Expand Up @@ -129,7 +128,7 @@ func NewMetaTemplater(templates cfg.MetaTemplates, logger *zap.Logger, cacheSize
}
}

cache, err := lru.New[string, MetaData](cacheSize)
cache, err := lru.New[uint64, MetaData](cacheSize)
if err != nil {
panic(err)
}
Expand All @@ -150,14 +149,15 @@ func NewMetaTemplater(templates cfg.MetaTemplates, logger *zap.Logger, cacheSize

type Data interface {
GetData() map[string]any
GetCacheKey() uint64
}

func (m *MetaTemplater) Render(data Data) (MetaData, error) {
initValues := data.GetData()
meta := MetaData{}

// Create a unique cache key based on the input data
cacheKey := generateCacheKey(initValues)
cacheKey := data.GetCacheKey()

// Check if the result is already cached
if cachedMeta, found := m.cache.Get(cacheKey); found {
Expand Down Expand Up @@ -211,36 +211,3 @@ func (m *MetaTemplater) Render(data Data) (MetaData, error) {

return meta, nil
}

func generateCacheKey(data map[string]any) string {
var builder strings.Builder
builder.Grow(len(data) * 16) // Preallocate memory for the builder (estimate)

for k, v := range data {
switch v := v.(type) {
case string:
// Write the key and string value to the builder
builder.WriteString(k)
builder.WriteString(":")
builder.WriteString(v)
builder.WriteString("|")
case int:
// Write the key and integer value to the builder
builder.WriteString(k)
builder.WriteString(":")
builder.WriteString(strconv.Itoa(v))
builder.WriteString("|")
}
// If the value is not a string or int, skip it
}

// Convert the builder to a string
key := builder.String()

// Remove the last "|" character if needed
if key != "" {
key = key[:len(key)-1] // Slice to remove the last character
}

return key
}
74 changes: 74 additions & 0 deletions pipeline/metadata/templater_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package metadata

import (
"fmt"
"strconv"
"testing"

"github.com/ozontech/file.d/cfg"
Expand Down Expand Up @@ -142,6 +143,24 @@ func TestTemplaterRender(t *testing.T) {
"broker": "kafka1:9093",
},
},
{
name: "Cache key includes int32/int64 values (kafka-like)",
templates: cfg.MetaTemplates{
"topic": "{{ .topic }}",
"partition": "partition_{{ .partition }}",
"offset": "offset_{{ .offset }}",
},
data: map[string]any{
"topic": "topic",
"partition": int32(1),
"offset": int64(100),
},
expected: map[string]any{
"topic": "topic",
"partition": "partition_1",
"offset": "offset_100",
},
},
}

for _, tt := range tests {
Expand All @@ -160,6 +179,57 @@ func TestTemplaterRender(t *testing.T) {
}
}

func TestMetaTemplaterCacheKey(t *testing.T) {
templater := NewMetaTemplater(
cfg.MetaTemplates{
"topic": "{{ .topic }}",
"partition": "partition_{{ .partition }}",
"offset": "offset_{{ .offset }}",
},
zap.NewExample(),
32,
)

// first record with partition=1, offset=100
first, err := templater.Render(metaInfoTest{
topic: "topic", partition: int32(1), offset: int64(100),
})
assert.Nil(t, err)

// first record with partition=2, offset=200
second, err := templater.Render(metaInfoTest{
topic: "topic", partition: int32(2), offset: int64(200),
})
assert.Nil(t, err)

assert.Equal(t, "partition_1", first["partition"])
assert.Equal(t, "offset_100", first["offset"])
assert.Equal(t, "partition_2", second["partition"])
assert.Equal(t, "offset_200", second["offset"])
}

type metaInfoTest struct {
topic string
partition int32
offset int64
}

func (m metaInfoTest) GetData() map[string]any {
return map[string]any{
"topic": m.topic,
"partition": m.partition,
"offset": m.offset,
}
}

func (m metaInfoTest) GetCacheKey() uint64 {
return Hash(
m.topic,
strconv.FormatInt(int64(m.partition), 10),
strconv.FormatInt(m.offset, 10),
)
}

type testMetadata struct {
data map[string]any
}
Expand All @@ -168,6 +238,10 @@ func (f testMetadata) GetData() map[string]any {
return f.data
}

func (f testMetadata) GetCacheKey() uint64 {
return Hash(fmt.Sprint(f.data))
}

func BenchmarkMetaTemplater_Render(b *testing.B) {
templater := NewMetaTemplater(
cfg.MetaTemplates{
Expand Down
20 changes: 20 additions & 0 deletions plugin/input/file/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,9 @@ import (
"os"
"os/exec"
"path/filepath"
"strconv"
"strings"
"time"

"github.com/ozontech/file.d/pipeline"
"github.com/ozontech/file.d/pipeline/metadata"
Expand Down Expand Up @@ -333,6 +335,24 @@ func (m metaInformation) GetData() map[string]any {
return data
}

func (m metaInformation) GetCacheKey() uint64 {
if m.k8sMetadata != nil {
return metadata.Hash(
m.k8sMetadata.PodName,
m.k8sMetadata.Namespace,
m.k8sMetadata.ContainerName,
m.k8sMetadata.GetPodStatus(),
m.k8sMetadata.GetUpdateTime().Format(time.RFC3339Nano),
)
}

return metadata.Hash(
m.filename,
m.symlink,
strconv.FormatUint(m.inode, 10),
)
}

/*{ meta-params
**`filename`**

Expand Down
1 change: 1 addition & 0 deletions plugin/input/file/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -396,6 +396,7 @@ func TestNewMetaInformation(t *testing.T) {
assert.Equal(t, tt.filename, metaInfo.filename)
assert.Equal(t, tt.symlink, metaInfo.symlink)
assert.Equal(t, uint64(tt.inode), metaInfo.inode)
assert.NotNil(t, metaInfo.GetCacheKey())

if tt.parseK8sMeta {
assert.Equal(t, tt.expectedK8sMeta.PodName, metaInfo.k8sMetadata.PodName)
Expand Down
4 changes: 3 additions & 1 deletion plugin/input/http/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,5 +142,7 @@ Key uses in the http_input_total metric.

**`params`** *`url.Values`*

**`request_uuid`** *`string`*
**`headers`** *`http.Header`*

**`request_uuid`** *`uint64`*
<br>*Generated using [__insane-doc__](https://github.com/vitkovskii/insane-doc)*
41 changes: 19 additions & 22 deletions plugin/input/http/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,16 @@ package http

import (
"context"
"crypto/sha1"
"fmt"
"io"
"net"
"net/http"
"net/url"
"strconv"
"strings"
"sync"
"time"

"github.com/google/uuid"
"github.com/klauspost/compress/gzip"
"github.com/ozontech/file.d/cfg"
"github.com/ozontech/file.d/fd"
Expand Down Expand Up @@ -415,7 +414,8 @@ func (p *Plugin) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if !ok {
p.failedAuthTotal.Inc()
p.errorsTotal.Inc()
p.logger.Warn("auth failed",
p.logger.Warn(
"auth failed",
zap.String("user_agent", r.UserAgent()),
zap.Any("headers", r.Header),
zap.String("remote_addr", r.RemoteAddr),
Expand Down Expand Up @@ -696,34 +696,29 @@ func newMetaInformation(login string, ip net.IP, r *http.Request) metaInformatio
}

func (m metaInformation) GetData() map[string]any {
contentLength := fmt.Sprintf("%d", m.request.ContentLength)
encodedParams := m.params.Encode()
remoteAddress := m.remoteAddr
result := fmt.Sprintf("%s|%s|%s", contentLength, encodedParams, remoteAddress)
requestUuid, _ := stringToUUID(result)

return map[string]any{
"login": m.login,
"remote_addr": m.remoteAddr,
"request": m.request,
"params": m.params,
"request_uuid": requestUuid.String(),
"headers": m.request.Header,
"request_uuid": m.cacheKey(),
}
}

func stringToUUID(input string) (uuid.UUID, error) {
hash := sha1.New()
_, err := hash.Write([]byte(input))
if err != nil {
return uuid.UUID{}, err
}

hashBytes := hash.Sum(nil)
func (m metaInformation) GetCacheKey() uint64 {
return m.cacheKey()
}

var u uuid.UUID
copy(u[:], hashBytes[:16])
func (m metaInformation) cacheKey() uint64 {
contentLength := strconv.FormatInt(m.request.ContentLength, 10)

return u, nil
return metadata.Hash(
contentLength,
m.params.Encode(),
m.remoteAddr.String(),
url.Values(m.request.Header).Encode(),
)
}

/*{ meta-params
Expand All @@ -735,5 +730,7 @@ func stringToUUID(input string) (uuid.UUID, error) {

**`params`** *`url.Values`*

**`request_uuid`** *`string`*
**`headers`** *`http.Header`*

**`request_uuid`** *`uint64`*
}*/
Loading
Loading