diff --git a/pkg/workflows/dontime/pb/dontime.pb.go b/pkg/workflows/dontime/pb/dontime.pb.go index 36d662710f..81a2403c5a 100644 --- a/pkg/workflows/dontime/pb/dontime.pb.go +++ b/pkg/workflows/dontime/pb/dontime.pb.go @@ -22,12 +22,11 @@ const ( ) type Observation struct { - state protoimpl.MessageState `protogen:"open.v1"` - Timestamp int64 `protobuf:"varint,1,opt,name=timestamp,proto3" json:"timestamp,omitempty"` - Requests map[string]int64 `protobuf:"bytes,2,rep,name=requests,proto3" json:"requests,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"` - LimitByBatchSizeFlag bool `protobuf:"varint,4,opt,name=limit_by_batch_size_flag,json=limitByBatchSizeFlag,proto3" json:"limit_by_batch_size_flag,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + state protoimpl.MessageState `protogen:"open.v1"` + Timestamp int64 `protobuf:"varint,1,opt,name=timestamp,proto3" json:"timestamp,omitempty"` + Requests map[string]int64 `protobuf:"bytes,2,rep,name=requests,proto3" json:"requests,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *Observation) Reset() { @@ -74,13 +73,6 @@ func (x *Observation) GetRequests() map[string]int64 { return nil } -func (x *Observation) GetLimitByBatchSizeFlag() bool { - if x != nil { - return x.LimitByBatchSizeFlag - } - return false -} - type ObservedDonTimes struct { state protoimpl.MessageState `protogen:"open.v1"` Timestamps []int64 `protobuf:"varint,1,rep,packed,name=timestamps,proto3" json:"timestamps,omitempty"` @@ -181,14 +173,13 @@ var File_dontime_proto protoreflect.FileDescriptor const file_dontime_proto_rawDesc = "" + "\n" + - "\rdontime.proto\"\xde\x01\n" + + "\rdontime.proto\"\xac\x01\n" + "\vObservation\x12\x1c\n" + "\ttimestamp\x18\x01 \x01(\x03R\ttimestamp\x126\n" + - "\brequests\x18\x02 \x03(\v2\x1a.Observation.RequestsEntryR\brequests\x126\n" + - "\x18limit_by_batch_size_flag\x18\x04 \x01(\bR\x14limitByBatchSizeFlag\x1a;\n" + + "\brequests\x18\x02 \x03(\v2\x1a.Observation.RequestsEntryR\brequests\x1a;\n" + "\rRequestsEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01J\x04\b\x03\x10\x04\"2\n" + + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01J\x04\b\x03\x10\x04J\x04\b\x04\x10\x05\"2\n" + "\x10ObservedDonTimes\x12\x1e\n" + "\n" + "timestamps\x18\x01 \x03(\x03R\n" + diff --git a/pkg/workflows/dontime/pb/dontime.proto b/pkg/workflows/dontime/pb/dontime.proto index d0baecf3b3..6c2d830aac 100644 --- a/pkg/workflows/dontime/pb/dontime.proto +++ b/pkg/workflows/dontime/pb/dontime.proto @@ -5,8 +5,7 @@ option go_package = "github.com/smartcontractkit/chainlink-common/pkg/workflows/ message Observation { int64 timestamp = 1; map requests = 2; - reserved 3; - bool limit_by_batch_size_flag = 4; + reserved 3, 4; } message ObservedDonTimes { diff --git a/pkg/workflows/dontime/plugin.go b/pkg/workflows/dontime/plugin.go index 5c91cbbfb9..24e79d1e35 100644 --- a/pkg/workflows/dontime/plugin.go +++ b/pkg/workflows/dontime/plugin.go @@ -196,9 +196,8 @@ func (p *Plugin) Observation(ctx context.Context, outctx ocr3types.OutcomeContex p.metrics.observationBatchOverflow.Record(ctx, int64(overflowCount)) observation := &pb.Observation{ - Timestamp: time.Now().UTC().UnixMilli(), - Requests: requests, - LimitByBatchSizeFlag: true, + Timestamp: time.Now().UTC().UnixMilli(), + Requests: requests, } return proto.MarshalOptions{Deterministic: true}.Marshal(observation) diff --git a/pkg/workflows/dontime/plugin_test.go b/pkg/workflows/dontime/plugin_test.go index 5ec9a139e2..8d81184fd1 100644 --- a/pkg/workflows/dontime/plugin_test.go +++ b/pkg/workflows/dontime/plugin_test.go @@ -561,13 +561,12 @@ func TestPlugin_Outcome_TrimByBatchSize(t *testing.T) { require.NoError(t, err) timestamp := time.Now().UnixMilli() - makeObservations := func(limitByBatchSize bool) []types.AttributedObservation { + makeObservations := func() []types.AttributedObservation { aos := make([]types.AttributedObservation, 4) for i := range 4 { obs := &pb.Observation{ - Timestamp: timestamp + int64(i), - Requests: map[string]int64{}, - LimitByBatchSizeFlag: limitByBatchSize, + Timestamp: timestamp + int64(i), + Requests: map[string]int64{}, } rawObs, err := proto.Marshal(obs) require.NoError(t, err) @@ -590,8 +589,8 @@ func TestPlugin_Outcome_TrimByBatchSize(t *testing.T) { prevOutcomeBytes, err := proto.Marshal(prevOutcome) require.NoError(t, err) - t.Run("trims when all observations set batch size flag", func(t *testing.T) { - outcome, err := plugin.Outcome(ctx, ocr3types.OutcomeContext{PreviousOutcome: prevOutcomeBytes}, query, makeObservations(true)) + t.Run("trims observed don times exceeding batch size", func(t *testing.T) { + outcome, err := plugin.Outcome(ctx, ocr3types.OutcomeContext{PreviousOutcome: prevOutcomeBytes}, query, makeObservations()) require.NoError(t, err) outcomeProto := &pb.Outcome{}