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
19 changes: 8 additions & 11 deletions fluent/fluent.go
Original file line number Diff line number Diff line change
Expand Up @@ -226,12 +226,12 @@ func newWithDialer(config Config, d dialer) (f *Fluent, err error) {
// "john smith",
// }
// f.Post("tag_name", structData)
func (f *Fluent) Post(tag string, message interface{}) error {
func (f *Fluent) Post(tag string, message any) error {
timeNow := time.Now()
return f.PostWithTime(tag, timeNow, message)
}

func (f *Fluent) PostWithTime(tag string, tm time.Time, message interface{}) error {
func (f *Fluent) PostWithTime(tag string, tm time.Time, message any) error {
if len(f.TagPrefix) > 0 {
tag = f.TagPrefix + "." + tag
}
Expand All @@ -245,9 +245,9 @@ func (f *Fluent) PostWithTime(tag string, tm time.Time, message interface{}) err

if msgtype.Kind() == reflect.Struct {
// message should be tagged by "codec" or "msg"
kv := make(map[string]interface{})
kv := make(map[string]any)
fields := msgtype.NumField()
for i := 0; i < fields; i++ {
for i := range fields {
field := msgtype.Field(i)
value := msg.FieldByIndex(field.Index)
// ignore unexported fields
Expand All @@ -271,15 +271,15 @@ func (f *Fluent) PostWithTime(tag string, tm time.Time, message interface{}) err
return errors.New("fluent#PostWithTime: map keys must be strings")
}

kv := make(map[string]interface{})
kv := make(map[string]any)
for _, k := range msg.MapKeys() {
kv[k.String()] = msg.MapIndex(k).Interface()
}

return f.EncodeAndPostData(tag, tm, kv)
}

func (f *Fluent) EncodeAndPostData(tag string, tm time.Time, message interface{}) error {
func (f *Fluent) EncodeAndPostData(tag string, tm time.Time, message any) error {
var msg *msgToSend
var err error
if msg, err = f.EncodeData(tag, tm, message); err != nil {
Expand Down Expand Up @@ -346,7 +346,7 @@ func getUniqueID(timeUnix int64) (string, error) {
return buf.String(), nil
}

func (f *Fluent) EncodeData(tag string, tm time.Time, message interface{}) (msg *msgToSend, err error) {
func (f *Fluent) EncodeData(tag string, tm time.Time, message any) (msg *msgToSend, err error) {
option := make(map[string]string)
msg = &msgToSend{}
timeUnix := tm.Unix()
Expand Down Expand Up @@ -536,10 +536,7 @@ func (f *Fluent) connectWithRetry(ctx context.Context) error {
return errIsClosing
}

waitTime := f.Config.RetryWait * e(defaultReconnectWaitIncreRate, float64(i-1))
if waitTime > f.Config.MaxRetryWait {
waitTime = f.Config.MaxRetryWait
}
waitTime := min(f.Config.RetryWait*e(defaultReconnectWaitIncreRate, float64(i-1)), f.Config.MaxRetryWait)

timeout = time.NewTimer(time.Duration(waitTime) * time.Millisecond)
case <-ctx.Done():
Expand Down
12 changes: 6 additions & 6 deletions fluent/fluent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ import (
"encoding/json"
"errors"
"fmt"
"io/ioutil"
"net"
"os"
"reflect"
"runtime"
"strconv"
Expand Down Expand Up @@ -166,7 +166,7 @@ func (d *testDialer) waitForNextDialing(accept bool, delayReads bool) *Conn {

// asserEqual asserts that actual and expected are equivalent, and otherwise
// marks the test as failed (t.Error). It uses reflect.DeepEqual internally.
func assertEqual(t *testing.T, actual, expected interface{}) {
func assertEqual(t *testing.T, actual, expected any) {
t.Helper()
if !reflect.DeepEqual(actual, expected) {
t.Errorf("got: '%+v', expected: '%+v'", actual, expected)
Expand Down Expand Up @@ -421,7 +421,7 @@ func Test_MarshalAsJSON(t *testing.T) {
}

func TestJsonConfig(t *testing.T) {
b, err := ioutil.ReadFile(`testdata/config.json`)
b, err := os.ReadFile(`testdata/config.json`)
if err != nil {
t.Error(err)
}
Expand Down Expand Up @@ -704,7 +704,7 @@ func TestNoPanicOnFailingTLSConnect(t *testing.T) {
}
defer f.Close()

for i := 0; i < 2; i++ {
for range 2 {
if err := f.Post("tag", map[string]string{"log": "msg"}); err == nil {
t.Error("Expected an error posting to an unreachable TLS endpoint")
}
Expand Down Expand Up @@ -830,10 +830,10 @@ func TestPendingChannelThreadSafety(t *testing.T) {
var wg sync.WaitGroup
wg.Add(numGoroutines)

for i := 0; i < numGoroutines; i++ {
for i := range numGoroutines {
go func(id int) {
defer wg.Done()
for j := 0; j < messagesPerGoroutine; j++ {
for j := range messagesPerGoroutine {
// Post a message
err := f.Post("tag", map[string]string{
"goroutine": strconv.Itoa(id),
Expand Down
18 changes: 9 additions & 9 deletions fluent/proto.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@ import (

//msgp:tuple Entry
type Entry struct {
Time int64 `msg:"time"`
Record interface{} `msg:"record"`
Time int64 `msg:"time"`
Record any `msg:"record"`
}

//msgp:tuple Forward
Expand All @@ -24,17 +24,17 @@ type Forward struct {

//msgp:tuple Message
type Message struct {
Tag string `msg:"tag"`
Time int64 `msg:"time"`
Record interface{} `msg:"record"`
Tag string `msg:"tag"`
Time int64 `msg:"time"`
Record any `msg:"record"`
Option map[string]string
}

//msgp:tuple MessageExt
type MessageExt struct {
Tag string `msg:"tag"`
Time EventTime `msg:"time,extension"`
Record interface{} `msg:"record"`
Tag string `msg:"tag"`
Time EventTime `msg:"time,extension"`
Record any `msg:"record"`
Option map[string]string
}

Expand Down Expand Up @@ -98,7 +98,7 @@ func (t *EventTime) MarshalBinaryTo(b []byte) error {
// UnmarshalBinary is implemented for testing and general completeness.
func (t *EventTime) UnmarshalBinary(b []byte) error {
if len(b) != length {
return fmt.Errorf("Invalid EventTime byte length: %d", len(b))
return fmt.Errorf("invalid EventTime byte length: %d", len(b))
}

sec := (int32(b[0]) << 24) | (int32(b[1]) << 16)
Expand Down
Loading