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
41 changes: 41 additions & 0 deletions encoder/encoder.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
package encoder

import (
"encoding/json"
"fmt"

"github.com/ozontech/file.d/pipeline"
)

const (
EncoderTypeJSON = "json"
EncoderTypeRaw = "raw"
)

type Encoder interface {
Encode(event *pipeline.Event, buf []byte) ([]byte, error)
}

type EncodingConfig struct {
Type string `json:"type" default:"json" options:"json|raw"`
Params json.RawMessage `json:"params"`
}

func NewEncoder(cfg EncodingConfig) (Encoder, error) {
switch cfg.Type {
case EncoderTypeJSON, "":
return newJSONEncoder(&JSONEncoderParams{}), nil

case EncoderTypeRaw:
var params RawEncoderParams
if len(cfg.Params) > 0 {
if err := json.Unmarshal(cfg.Params, &params); err != nil {
return nil, fmt.Errorf("raw encoder params: %w", err)
}
}
return newRawEncoder(&params), nil

default:
return nil, fmt.Errorf("unknown encoding type %q; supported: json, raw", cfg.Type)
}
}
153 changes: 153 additions & 0 deletions encoder/encoder_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
package encoder

import (
"encoding/json"
"testing"

"github.com/ozontech/file.d/pipeline"
insaneJSON "github.com/ozontech/insane-json"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func newTestEvent(t testing.TB, raw string) *pipeline.Event {
t.Helper()

root, err := insaneJSON.DecodeString(raw)
require.NoError(t, err)

t.Cleanup(func() {
insaneJSON.Release(root)
})

return &pipeline.Event{Root: root}
}

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

tests := []struct {
name string
cfg EncodingConfig
wantErr bool
errContains string
assertType func(t *testing.T, enc Encoder)
}{
{
name: "explicit json type",
cfg: EncodingConfig{Type: EncoderTypeJSON},
assertType: func(t *testing.T, enc Encoder) {
assert.IsType(t, &JSONEncoder{}, enc)
},
},
{
name: "empty type defaults to json",
cfg: EncodingConfig{Type: ""},
assertType: func(t *testing.T, enc Encoder) {
assert.IsType(t, &JSONEncoder{}, enc)
},
},
{
name: "raw type without params uses default field",
cfg: EncodingConfig{Type: EncoderTypeRaw},
assertType: func(t *testing.T, enc Encoder) {
raw, ok := enc.(*RawEncoder)
require.True(t, ok)
assert.Equal(t, "message", raw.field)
},
},
{
name: "raw type with empty params object uses default field",
cfg: EncodingConfig{Type: EncoderTypeRaw, Params: json.RawMessage(`{}`)},
assertType: func(t *testing.T, enc Encoder) {
raw, ok := enc.(*RawEncoder)
require.True(t, ok)
assert.Equal(t, "message", raw.field)
},
},
{
name: "raw type with custom field",
cfg: EncodingConfig{Type: EncoderTypeRaw, Params: json.RawMessage(`{"field":"data"}`)},
assertType: func(t *testing.T, enc Encoder) {
raw, ok := enc.(*RawEncoder)
require.True(t, ok)
assert.Equal(t, "data", raw.field)
},
},
{
name: "raw type with empty field falls back to message",
cfg: EncodingConfig{Type: EncoderTypeRaw, Params: json.RawMessage(`{"field":""}`)},
assertType: func(t *testing.T, enc Encoder) {
raw, ok := enc.(*RawEncoder)
require.True(t, ok)
assert.Equal(t, "message", raw.field)
},
},
{
name: "raw type with invalid params",
cfg: EncodingConfig{Type: EncoderTypeRaw, Params: json.RawMessage(`{"field":`)},
wantErr: true,
errContains: "raw encoder params",
},
{
name: "unknown type",
cfg: EncodingConfig{Type: "yaml"},
wantErr: true,
errContains: `unknown encoding type "yaml"`,
},
}

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

enc, err := NewEncoder(tt.cfg)

if tt.wantErr {
require.Error(t, err)
assert.Nil(t, enc)
if tt.errContains != "" {
assert.Contains(t, err.Error(), tt.errContains)
}
return
}

require.NoError(t, err)
require.NotNil(t, enc)
if tt.assertType != nil {
tt.assertType(t, enc)
}
})
}
}

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

t.Run("json", func(t *testing.T) {
t.Parallel()

enc, err := NewEncoder(EncodingConfig{Type: EncoderTypeJSON})
require.NoError(t, err)

event := newTestEvent(t, `{"message":"hi"}`)
out, err := enc.Encode(event, nil)
require.NoError(t, err)
assert.JSONEq(t, `{"message":"hi"}`, string(out))
})

t.Run("raw", func(t *testing.T) {
t.Parallel()

enc, err := NewEncoder(EncodingConfig{
Type: EncoderTypeRaw,
Params: json.RawMessage(`{"field":"message"}`),
})
require.NoError(t, err)

event := newTestEvent(t, `{"message":"hi"}`)
out, err := enc.Encode(event, nil)
require.NoError(t, err)
assert.Equal(t, "hi", string(out))
})
}
16 changes: 16 additions & 0 deletions encoder/json.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
package encoder

import "github.com/ozontech/file.d/pipeline"

type JSONEncoderParams struct{}

type JSONEncoder struct{}

func newJSONEncoder(_ *JSONEncoderParams) *JSONEncoder {
return &JSONEncoder{}
}

func (e *JSONEncoder) Encode(event *pipeline.Event, buf []byte) ([]byte, error) {
buf, _ = event.Encode(buf)
return buf, nil
}
67 changes: 67 additions & 0 deletions encoder/json_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
package encoder

import (
"testing"

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

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

tests := []struct {
name string
input string
expected string
}{
{
name: "simple object",
input: `{"message":"hello"}`,
expected: `{"message":"hello"}`,
},
{
name: "nested object",
input: `{"level":"info","data":{"a":1,"b":[1,2,3]}}`,
expected: `{"level":"info","data":{"a":1,"b":[1,2,3]}}`,
},
{
name: "array root",
input: `[1,2,3]`,
expected: `[1,2,3]`,
},
{
name: "scalar string root",
input: `"just a string"`,
expected: `"just a string"`,
},
{
name: "number root",
input: `42`,
expected: `42`,
},
{
name: "empty object",
input: `{}`,
expected: `{}`,
},
{
name: "special characters escaped",
input: `{"msg":"line1\nline2\t\"quoted\""}`,
expected: `{"msg":"line1\nline2\t\"quoted\""}`,
},
}

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

enc := newJSONEncoder(&JSONEncoderParams{})
event := newTestEvent(t, tt.input)

out, err := enc.Encode(event, nil)
require.NoError(t, err)
assert.JSONEq(t, tt.expected, string(out))
})
}
}
41 changes: 41 additions & 0 deletions encoder/raw.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
package encoder

import (
"errors"
"fmt"

"github.com/ozontech/file.d/pipeline"
)

const defaultRawField = "message"

var ErrFieldNotFound = errors.New("field not found")

type RawEncoderParams struct {
Field string `json:"field"`
}

type RawEncoder struct {
field string
}

func newRawEncoder(params *RawEncoderParams) *RawEncoder {
field := params.Field
if field == "" {
field = defaultRawField
}
return &RawEncoder{field: field}
}

func (e *RawEncoder) Encode(event *pipeline.Event, buf []byte) ([]byte, error) {
node := event.Root.Dig(e.field)
if node == nil {
return buf, fmt.Errorf("%w: %q", ErrFieldNotFound, e.field)
}

if node.IsString() {
return append(buf, node.AsBytes()...), nil
}

return node.Encode(buf), nil
}
Loading
Loading