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
45 changes: 43 additions & 2 deletions pkg/chipingress/batch/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,11 +14,18 @@ import (
"go.opentelemetry.io/otel/attribute"
otelmetric "go.opentelemetry.io/otel/metric"
"go.uber.org/zap"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"

"github.com/smartcontractkit/chainlink-common/pkg/chipingress"
)

var (
ErrMessageBufferFull = errors.New("message buffer is full")
ErrClientShutdown = errors.New("client is shutdown")
)

type messageWithCallback struct {
event *chipingress.CloudEventPb
callback func(error)
Expand Down Expand Up @@ -51,6 +58,40 @@ const (
ErrCodeResultsMismatch = chipingress.PublishErrorCode(-1) // client-side synthetic code
)

// ErrorCodeFor returns a bounded error code string for a batch send/queue
// error, suitable for use as a metric dimension.
func ErrorCodeFor(err error) string {
if err == nil {
return ""
}

var pubErr *PublishError
if errors.As(err, &pubErr) {
if pubErr.Code == ErrCodeResultsMismatch {
return "results_mismatch"
}
return pubErr.Code.String()
}

if errors.Is(err, ErrMessageBufferFull) {
return ErrMessageBufferFull.Error()
}
if errors.Is(err, ErrClientShutdown) {
return ErrClientShutdown.Error()
}
Comment on lines +76 to +81
if st, ok := status.FromError(err); ok {
return st.Code().String()
}
Comment on lines +82 to +84
if errors.Is(err, context.DeadlineExceeded) {
return codes.DeadlineExceeded.String()
}
if errors.Is(err, context.Canceled) {
return codes.Canceled.String()
}

return "client_error"
}

// Client is a batching client that accumulates messages and sends them in batches.
type Client struct {
client chipingress.Client
Expand Down Expand Up @@ -242,7 +283,7 @@ func (b *Client) QueueMessage(event *chipingress.CloudEventPb, callback func(err
// Check shutdown first to avoid race with buffer send
select {
case <-b.stopCh:
return errors.New("client is shutdown")
return ErrClientShutdown
default:
}

Expand Down Expand Up @@ -276,7 +317,7 @@ func (b *Client) QueueMessage(event *chipingress.CloudEventPb, callback func(err
case b.messageBuffer <- msg:
return nil
default:
return errors.New("message buffer is full")
return ErrMessageBufferFull
}
}

Expand Down
70 changes: 70 additions & 0 deletions pkg/chipingress/batch/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@ import (
"go.opentelemetry.io/otel/attribute"
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
"go.opentelemetry.io/otel/sdk/metric/metricdata"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"

"github.com/smartcontractkit/chainlink-common/pkg/chipingress"
Expand Down Expand Up @@ -2346,3 +2348,71 @@ func TestTransactionEnabledEdgeCases(t *testing.T) {
assert.EqualError(t, err3, "kafka unavailable")
})
}

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

tests := []struct {
name string
err error
want string
}{
{
name: "partial delivery publish error",
err: &PublishError{Code: chipingress.PublishErrorCode(1), Reason: "schema not found"},
want: chipingress.PublishErrorCode(1).String(),
},
{
name: "results mismatch",
err: &PublishError{Code: ErrCodeResultsMismatch, Reason: "server returned 1 results for 2 events"},
want: "results_mismatch",
},
{
name: "deadline exceeded status",
err: status.Error(codes.DeadlineExceeded, "context deadline exceeded"),
want: codes.DeadlineExceeded.String(),
},
{
name: "deadline exceeded context",
err: context.DeadlineExceeded,
want: codes.DeadlineExceeded.String(),
},
{
name: "unavailable gateway 502",
err: status.Error(codes.Unavailable, `unexpected HTTP status code received from server: 502 (Bad Gateway)`),
want: codes.Unavailable.String(),
},
{
name: "internal publish failure",
err: status.Error(codes.Internal, "failed to publish events"),
want: codes.Internal.String(),
},
{
name: "buffer full",
err: ErrMessageBufferFull,
want: ErrMessageBufferFull.Error(),
},
{
name: "client shutdown",
err: ErrClientShutdown,
want: ErrClientShutdown.Error(),
},
{
name: "unknown error",
err: errors.New("something else"),
want: "client_error",
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
assert.Equal(t, tt.want, ErrorCodeFor(tt.err))
})
}
}

func TestErrorCodeFor_nil(t *testing.T) {
t.Parallel()
require.Empty(t, ErrorCodeFor(nil))
}
Loading