diff --git a/consensus/oracle/plugin/plugin.go b/consensus/oracle/plugin/plugin.go index 7fa596f7b..86e384a8d 100644 --- a/consensus/oracle/plugin/plugin.go +++ b/consensus/oracle/plugin/plugin.go @@ -77,10 +77,6 @@ func ToRequestMetaData(metadata oracle.ConsensusRequestMetadata) *oracletypes.Re } } -func (r *reportingPlugin) ValidateObservation(ctx context.Context, outctx ocr3types.OutcomeContext, query types.Query, ao types.AttributedObservation) error { - return nil -} - func (r *reportingPlugin) ObservationQuorum(ctx context.Context, outctx ocr3types.OutcomeContext, query types.Query, aos []types.AttributedObservation) (bool, error) { return quorumhelper.ObservationCountReachesObservationQuorum(quorumhelper.QuorumTwoFPlusOne, r.n, r.f, aos), nil } diff --git a/consensus/oracle/plugin/plugin_outcome.go b/consensus/oracle/plugin/plugin_outcome.go index a5b3111a0..a8a5a786f 100644 --- a/consensus/oracle/plugin/plugin_outcome.go +++ b/consensus/oracle/plugin/plugin_outcome.go @@ -28,6 +28,30 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/logger" ) +func (r *reportingPlugin) ValidateObservation(ctx context.Context, outctx ocr3types.OutcomeContext, query types.Query, ao types.AttributedObservation) error { + lggr := logger.With(r.lggr, "seqNr", outctx.SeqNr, "observer", ao.Observer) + + obs := &oracletypes.Observation{} + if err := proto.Unmarshal(ao.Observation, obs); err != nil { + lggr.Warnw("could not unmarshal observation", "error", err) + return fmt.Errorf("could not unmarshal observation from observer %d: %w", ao.Observer, err) + } + + for requestID, reqObs := range obs.Observations { + if reqObs.Metadata == nil { + lggr.Warnw("observation missing metadata", "requestID", requestID) + return fmt.Errorf("observation from observer %d is missing metadata for request %s", ao.Observer, requestID) + } + + if reqObs.Input == nil { + lggr.Warnw("observation missing input", "requestID", requestID) + return fmt.Errorf("observation from observer %d is missing input for request %s", ao.Observer, requestID) + } + } + + return nil +} + func (r *reportingPlugin) Outcome(ctx context.Context, outctx ocr3types.OutcomeContext, query types.Query, attributedObservations []types.AttributedObservation) (ocr3types.Outcome, error) { lggr := logger.With(r.lggr, "seqNr", outctx.SeqNr) @@ -96,6 +120,12 @@ func (r *reportingPlugin) addRequestOutcomeToBatch(ctx context.Context, lggr log updateErrorHandlingFlag = false } + // Does the observation have a valid input? + if obs.Input == nil { + lggr.Warnw("observation missing input", "requestID", requestID, "observerMetadata", obs.Metadata) + continue + } + // Does the observation have a timestamp? if obs.ReceivedAt == nil { lggr.Warnw("observation missing receivedAt timestamp", "requestID", requestID, "observerMetadata", obs.Metadata) @@ -275,6 +305,10 @@ func (r *reportingPlugin) calculateConsensusMetadataDescriptorAndDefault(lggr lo updateErrorHandlingFlag bool) (*oracletypes.RequestObservation, error) { var allObservationsMDDBytes []*valuespb.Value for _, obs := range observations { + if obs.Input == nil { + lggr.Errorw("observation missing input, skipping MDD calculation", "observerMetadata", obs.Metadata) + continue + } mddBytes, err := proto.MarshalOptions{Deterministic: true}.Marshal(&oracletypes.RequestObservation{ Metadata: obs.Metadata, Input: &sdk.SimpleConsensusInputs{ diff --git a/consensus/oracle/plugin/plugin_outcome_test.go b/consensus/oracle/plugin/plugin_outcome_test.go index 555d6a0c9..07378f360 100644 --- a/consensus/oracle/plugin/plugin_outcome_test.go +++ b/consensus/oracle/plugin/plugin_outcome_test.go @@ -336,6 +336,103 @@ func Test_Outcome_IdenticalConsensus_failureCodes(t *testing.T) { }) } +func makeTestObs( + t *testing.T, + reqID string, + md oracle.ConsensusRequestMetadata, + observerID uint8, + input *sdk.SimpleConsensusInputs, +) libocrtypes.AttributedObservation { + t.Helper() + + ro := &oracletypes.RequestObservation{ + Metadata: plugin.ToRequestMetaData(md), + Input: input, + ReceivedAt: timestamppb.New(time.Now()), + } + + obsProto := &oracletypes.Observation{ + Observations: map[string]*oracletypes.RequestObservation{reqID: ro}, + } + b, err := proto.Marshal(obsProto) + require.NoError(t, err) + + return libocrtypes.AttributedObservation{ + Observation: b, + Observer: commontypes.OracleID(observerID), + } +} + +func Test_Outcome_NilInputs(t *testing.T) { + t.Parallel() + + lggr := logger.Test(t) + ctx := context.Background() + + const testF, testN = 2, 7 + reportingPlugin, _ := createReportingPlugin(t, lggr, testF, testN, 5, defaultMaxLengthBytes) + + md := testMetaData() + reqID := md.RequestID() + + // 2f+1 = 5 observations: one malformed (nil Input) + four valid values with MEDIAN + // aggregation. The four valid values [10,20,30,40] have median 25, so consensus + // should succeed. + attributed := []libocrtypes.AttributedObservation{ + makeTestObs(t, reqID, md, 0, nil), + makeOutcomeTestObs(t, reqID, md, sdk.AggregationType_AGGREGATION_TYPE_MEDIAN, 1, false, true, true), + makeOutcomeTestObs(t, reqID, md, sdk.AggregationType_AGGREGATION_TYPE_MEDIAN, 2, false, true, true), + makeOutcomeTestObs(t, reqID, md, sdk.AggregationType_AGGREGATION_TYPE_MEDIAN, 3, false, true, true), + makeOutcomeTestObs(t, reqID, md, sdk.AggregationType_AGGREGATION_TYPE_MEDIAN, 4, false, true, true), + } + + qBytes, err := proto.Marshal(&oracletypes.Query{RequestIDs: []string{reqID}}) + require.NoError(t, err) + + outcomeBytes, err := reportingPlugin.Outcome(ctx, ocr3types.OutcomeContext{SeqNr: 1}, qBytes, attributed) + require.NoError(t, err) + + // Consensus should succeed with the 4 valid identical observations. + outcome := &oracletypes.Outcome{} + require.NoError(t, proto.Unmarshal(outcomeBytes, outcome)) + require.Len(t, outcome.Outcomes, 1) + require.NotNil(t, outcome.Outcomes[0].GetSuccess(), "expected a successful consensus outcome") +} + +func Test_Outcome_AllNilInputs(t *testing.T) { + t.Parallel() + + lggr := logger.Test(t) + ctx := context.Background() + + const testF, testN = 2, 7 + reportingPlugin, _ := createReportingPlugin(t, lggr, testF, testN, 5, defaultMaxLengthBytes) + + md := testMetaData() + reqID := md.RequestID() + + // 2f+1 = 5 observations, all with nil Input. + attributed := []libocrtypes.AttributedObservation{ + makeTestObs(t, reqID, md, 0, nil), + makeTestObs(t, reqID, md, 1, nil), + makeTestObs(t, reqID, md, 2, nil), + makeTestObs(t, reqID, md, 3, nil), + makeTestObs(t, reqID, md, 4, nil), + } + + qBytes, err := proto.Marshal(&oracletypes.Query{RequestIDs: []string{reqID}}) + require.NoError(t, err) + + outcomeBytes, err := reportingPlugin.Outcome(ctx, ocr3types.OutcomeContext{SeqNr: 1}, qBytes, attributed) + require.NoError(t, err) + + // All observations are skipped, so consensus should fail. + outcome := &oracletypes.Outcome{} + require.NoError(t, proto.Unmarshal(outcomeBytes, outcome)) + require.Len(t, outcome.Outcomes, 1) + require.NotNil(t, outcome.Outcomes[0].GetFailure(), "expected a failed consensus outcome when all inputs are nil") +} + func Test_Outcome_RecordsObservationQuorumForTimeoutClassification(t *testing.T) { t.Parallel()