From 640ba8569acbf07931c4136f9d3d56d1033fcc47 Mon Sep 17 00:00:00 2001 From: ilija42 Date: Tue, 8 Sep 2026 22:24:14 +0200 Subject: [PATCH] Fail invalid consensus outcomes per request instead of signing them --- consensus/action/capability.go | 18 ++ consensus/action/capability_test.go | 9 + .../oracle/plugin/errors_consensus_test.go | 7 + consensus/oracle/plugin/plugin_observation.go | 1 + consensus/oracle/plugin/plugin_outcome.go | 13 + .../oracle/plugin/plugin_outcome_test.go | 112 ++++++++ consensus/oracle/plugin/plugin_reports.go | 188 +++++++++----- .../oracle/plugin/plugin_reports_test.go | 243 ++++++++++++++++++ .../oracle/types/value_consensus_types.pb.go | 23 +- .../oracle/types/value_consensus_types.proto | 4 + 10 files changed, 548 insertions(+), 70 deletions(-) create mode 100644 consensus/oracle/plugin/plugin_reports_test.go diff --git a/consensus/action/capability.go b/consensus/action/capability.go index 55d75eb16..b8642c592 100644 --- a/consensus/action/capability.go +++ b/consensus/action/capability.go @@ -284,6 +284,10 @@ func (c *consensusCapability) Simple(ctx context.Context, metadata capabilities. return nil, response.Err } + if err := validateRawReportHasPayload(response.RawReport); err != nil { + return nil, caperrors.NewPublicSystemError(fmt.Errorf("invalid report for request %s: %w", consensusRequestMetaData.RequestID(), err), caperrors.Internal) + } + // Remove the metadata prefix from the raw report to get the serialised value serialisedValue := response.RawReport[plugin.ReportMetaDataPrependLength:] @@ -376,6 +380,10 @@ func (c *consensusCapability) Report(ctx context.Context, metadata capabilities. return nil, response.Err } + if err := validateRawReportHasPayload(response.RawReport); err != nil { + return nil, caperrors.NewPublicSystemError(fmt.Errorf("invalid report for request %s: %w", consensusRequestMetaData.RequestID(), err), caperrors.Internal) + } + var sigs []*sdk.AttributedSignature for _, s := range response.Sigs { @@ -615,3 +623,13 @@ func decodeObservationType(lggr logger.Logger, input *sdk.SimpleConsensusInputs) } return nil } + +// validateRawReportHasPayload checks that a signed report carries a payload after the metadata prefix. Without this check +// a report consisting of only the prefix would be sliced into an empty payload, which unmarshals into a zero value that +// cannot be told apart from a genuine result. +func validateRawReportHasPayload(rawReport []byte) error { + if len(rawReport) <= plugin.ReportMetaDataPrependLength { + return fmt.Errorf("report of %d bytes has no payload after the %d byte metadata prefix", len(rawReport), plugin.ReportMetaDataPrependLength) + } + return nil +} diff --git a/consensus/action/capability_test.go b/consensus/action/capability_test.go index aea60d7b6..7b277cfba 100644 --- a/consensus/action/capability_test.go +++ b/consensus/action/capability_test.go @@ -23,6 +23,7 @@ import ( "github.com/smartcontractkit/chainlink-protos/cre/go/sdk" "github.com/smartcontractkit/chainlink-protos/cre/go/values" + "github.com/smartcontractkit/capabilities/consensus/oracle/plugin" "github.com/smartcontractkit/capabilities/libs/testutils" ) @@ -490,3 +491,11 @@ func generateRandomHexString(byteLength int) string { } return hex.EncodeToString(randomBytes) } + +func Test_validateRawReportHasPayload(t *testing.T) { + metadataPrefix := make([]byte, plugin.ReportMetaDataPrependLength) + + require.Error(t, validateRawReportHasPayload(nil)) + require.Error(t, validateRawReportHasPayload(metadataPrefix), "a report consisting of only the metadata prefix has no payload") + require.NoError(t, validateRawReportHasPayload(append(metadataPrefix, 0x01))) +} diff --git a/consensus/oracle/plugin/errors_consensus_test.go b/consensus/oracle/plugin/errors_consensus_test.go index e76ffb9ab..c8c8b2e2f 100644 --- a/consensus/oracle/plugin/errors_consensus_test.go +++ b/consensus/oracle/plugin/errors_consensus_test.go @@ -4,8 +4,10 @@ import ( "errors" "testing" + "github.com/stretchr/testify/require" "google.golang.org/protobuf/types/known/structpb" + ocrtypes "github.com/smartcontractkit/chainlink-common/pkg/capabilities/consensus/ocr3/types" "github.com/smartcontractkit/chainlink-common/pkg/logger" "github.com/smartcontractkit/capabilities/consensus/oracle" @@ -82,6 +84,11 @@ func Test_ReceivedTooManyErrorsWithDefault(t *testing.T) { newCrWithErrorAndDefault(t, errors.New("its broken"), 20, md1)}, verifyReport: func(t *testing.T, report ocr3types.ReportPlus[[]byte], infos *structpb.Struct) { verifyValueConsensusReport(t, report, infos, values.NewInt64(20), "evm") + + // The default value is timestamped with when the DON observed the request, not with the zero time + meta, _, err := ocrtypes.Decode(report.ReportWithInfo.Report) + require.NoError(t, err, "Failed to extract metadata fields from report") + require.NotZero(t, meta.Timestamp) }}, } diff --git a/consensus/oracle/plugin/plugin_observation.go b/consensus/oracle/plugin/plugin_observation.go index a6bc73061..1e3694240 100644 --- a/consensus/oracle/plugin/plugin_observation.go +++ b/consensus/oracle/plugin/plugin_observation.go @@ -36,6 +36,7 @@ func (r *reportingPlugin) Observation(ctx context.Context, outctx ocr3types.Outc Input: req.Input, RemoveLibUseInFailureMessageFormattingFlag: true, UpdateErrorHandlingFlag: true, + IncludeErrorObservationTimestampsFlag: true, } hasCapacity := observationBatch.AddObservation(ctx, reqObs) diff --git a/consensus/oracle/plugin/plugin_outcome.go b/consensus/oracle/plugin/plugin_outcome.go index a8a5a786f..46574e086 100644 --- a/consensus/oracle/plugin/plugin_outcome.go +++ b/consensus/oracle/plugin/plugin_outcome.go @@ -108,9 +108,11 @@ func (r *reportingPlugin) addRequestOutcomeToBatch(ctx context.Context, lggr log var obsErrors []string var obsValues []*valuespb.Value var timestamps []*timestamppb.Timestamp + var errorObsTimestamps []*timestamppb.Timestamp removeLibUseInErrorFormattingFlag := true updateErrorHandlingFlag := true + includeErrorObservationTimestampsFlag := true for _, obs := range observations { if !obs.RemoveLibUseInFailureMessageFormattingFlag { // enable only when all nodes are updated removeLibUseInErrorFormattingFlag = false @@ -120,6 +122,10 @@ func (r *reportingPlugin) addRequestOutcomeToBatch(ctx context.Context, lggr log updateErrorHandlingFlag = false } + if !obs.IncludeErrorObservationTimestampsFlag { + includeErrorObservationTimestampsFlag = false + } + // Does the observation have a valid input? if obs.Input == nil { lggr.Warnw("observation missing input", "requestID", requestID, "observerMetadata", obs.Metadata) @@ -145,9 +151,16 @@ func (r *reportingPlugin) addRequestOutcomeToBatch(ctx context.Context, lggr log timestamps = append(timestamps, obs.ReceivedAt) case *sdk.SimpleConsensusInputs_Error: obsErrors = append(obsErrors, inputObservation.Error) + errorObsTimestamps = append(errorObsTimestamps, obs.ReceivedAt) } } + // Error observations are included in the median timestamp so that an outcome which falls back to the default value + // after f+1 errors is timestamped with when the DON observed the request rather than with the zero time + if includeErrorObservationTimestampsFlag { + timestamps = append(timestamps, errorObsTimestamps...) + } + timestamp := ×tamppb.Timestamp{} if len(timestamps) > 0 { timestamp = calculateMedianTimestamp(timestamps) diff --git a/consensus/oracle/plugin/plugin_outcome_test.go b/consensus/oracle/plugin/plugin_outcome_test.go index 07378f360..f8a21a38c 100644 --- a/consensus/oracle/plugin/plugin_outcome_test.go +++ b/consensus/oracle/plugin/plugin_outcome_test.go @@ -126,6 +126,16 @@ func extractSingleFailureMessage(t *testing.T, outcomeBytes ocr3types.Outcome) s return failure.FailureMessage } +func extractSingleSuccessOutcome(t *testing.T, outcomeBytes ocr3types.Outcome) *oracletypes.ConsensusSuccessOutcome { + t.Helper() + outcome := &oracletypes.Outcome{} + require.NoError(t, proto.Unmarshal(outcomeBytes, outcome)) + require.Len(t, outcome.Outcomes, 1, "expected exactly one consensus outcome") + success := outcome.Outcomes[0].GetSuccess() + require.NotNil(t, success, "expected a successful consensus outcome, got a failure") + return success +} + func extractSingleFailureCode(t *testing.T, outcomeBytes ocr3types.Outcome) oracletypes.ConsensusFailureCode { t.Helper() outcome := &oracletypes.Outcome{} @@ -278,6 +288,108 @@ func Test_Outcome_RemoveLibUseInFailureMessageFormatting(t *testing.T) { }) } +// makeTimestampedOutcomeTestObs builds a single AttributedObservation with the given receivedAt time for direct Outcome() +// timestamp tests. Every observation carries a default value so that a request which receives f+1 errors still produces +// a successful outcome. +func makeTimestampedOutcomeTestObs( + t *testing.T, + reqID string, + md oracle.ConsensusRequestMetadata, + observerID uint8, + isError bool, + receivedAt time.Time, + includeErrorObservationTimestampsFlag bool, +) libocrtypes.AttributedObservation { + t.Helper() + + simpleInputs := &sdk.SimpleConsensusInputs{ + Descriptors: &sdk.ConsensusDescriptor{ + Descriptor_: &sdk.ConsensusDescriptor_Aggregation{Aggregation: sdk.AggregationType_AGGREGATION_TYPE_MEDIAN}, + }, + Default: values.Proto(values.NewInt64(20)), + } + if isError { + simpleInputs.Observation = &sdk.SimpleConsensusInputs_Error{Error: fmt.Sprintf("error from observer %d", observerID)} + } else { + simpleInputs.Observation = &sdk.SimpleConsensusInputs_Value{Value: values.Proto(values.NewInt64(int64(observerID) * 10))} + } + + ro := &oracletypes.RequestObservation{ + Metadata: plugin.ToRequestMetaData(md), + Input: simpleInputs, + ReceivedAt: timestamppb.New(receivedAt), + RemoveLibUseInFailureMessageFormattingFlag: true, + UpdateErrorHandlingFlag: true, + IncludeErrorObservationTimestampsFlag: includeErrorObservationTimestampsFlag, + } + + 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), + } +} + +// Test_Outcome_IncludeErrorObservationTimestamps documents how the outcome timestamp is calculated depending on +// RequestObservation.include_error_observation_timestamps_flag when f+1 errors are received and the default value is used. +func Test_Outcome_IncludeErrorObservationTimestamps(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() + + qBytes, err := proto.Marshal(&oracletypes.Query{RequestIDs: []string{reqID}}) + require.NoError(t, err) + + now := time.Now().Truncate(time.Second) + receivedAt := func(offset int) time.Time { + return now.Add(time.Duration(offset) * time.Second) + } + + // 2f+1 = 5 observations: f+1 errors received first followed by two values, so the outcome is the default value + newObservations := func(flag bool) []libocrtypes.AttributedObservation { + return []libocrtypes.AttributedObservation{ + makeTimestampedOutcomeTestObs(t, reqID, md, 0, true, receivedAt(1), flag), + makeTimestampedOutcomeTestObs(t, reqID, md, 1, true, receivedAt(2), flag), + makeTimestampedOutcomeTestObs(t, reqID, md, 2, true, receivedAt(3), flag), + makeTimestampedOutcomeTestObs(t, reqID, md, 3, false, receivedAt(4), flag), + makeTimestampedOutcomeTestObs(t, reqID, md, 4, false, receivedAt(5), flag), + } + } + + t.Run("flag_set_on_all_observations", func(t *testing.T) { + outcomeBytes, err := reportingPlugin.Outcome(ctx, ocr3types.OutcomeContext{SeqNr: 1}, qBytes, newObservations(true)) + require.NoError(t, err) + + // The timestamp is the median of all five observations + success := extractSingleSuccessOutcome(t, outcomeBytes) + assert.Equal(t, receivedAt(3).Unix(), success.Timestamp.AsTime().Unix()) + }) + + t.Run("flag_not_set_on_one_observation", func(t *testing.T) { + attributed := newObservations(true) + attributed[4] = makeTimestampedOutcomeTestObs(t, reqID, md, 4, false, receivedAt(5), false) + + outcomeBytes, err := reportingPlugin.Outcome(ctx, ocr3types.OutcomeContext{SeqNr: 1}, qBytes, attributed) + require.NoError(t, err) + + // The timestamp is the median of the two value observations only, as calculated by nodes without the flag + success := extractSingleSuccessOutcome(t, outcomeBytes) + assert.Equal(t, receivedAt(4).Unix(), success.Timestamp.AsTime().Unix()) + }) +} + // Test_Outcome_IdenticalConsensus_failureCodes documents how ConsensusFailureCode is chosen for // identical-consensus threshold failures depending on RequestObservation.update_error_handling_flag. func Test_Outcome_IdenticalConsensus_failureCodes(t *testing.T) { diff --git a/consensus/oracle/plugin/plugin_reports.go b/consensus/oracle/plugin/plugin_reports.go index 68c583711..8ba7ddd05 100644 --- a/consensus/oracle/plugin/plugin_reports.go +++ b/consensus/oracle/plugin/plugin_reports.go @@ -3,6 +3,7 @@ package plugin import ( "context" "fmt" + "math" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/structpb" @@ -34,67 +35,36 @@ func (r *reportingPlugin) Reports(ctx context.Context, seqNr uint64, outcome ocr // Create a report for each outcome for _, reqOutcome := range requestsOutcome.Outcomes { + // The limit is checked before the next outcome is processed so that every report added to the round, including + // failure reports, counts towards it + if r.maxNumberOfReports > 0 && len(reports) >= r.maxNumberOfReports { + r.lggr.Warnw("maximum number of reports reached, stopping further report generation for this round", "seqNr", seqNr, "maxNumberOfReports", r.maxNumberOfReports) + break + } + switch v := reqOutcome.GetOutcome().(type) { case *oracletypes.ConsensusOutcome_Success: successOutcome := v.Success - r.lggr.Debugw("received successful consensus outcome", "seqNr", seqNr, "requestID", successOutcome.Metadata.RequestId) reqMetadata := successOutcome.Metadata - var report []byte - switch reqMetadata.RequestType { - case oracletypes.RequestType_VALUE_CONSENSUS: - report = successOutcome.Outcome - case oracletypes.RequestType_REPORT_GENERATION: - // If the request type is report extract the report from the values.Value before signing it - serialisedValue := successOutcome.Outcome - value := &valuespb.Value{} - if err := proto.Unmarshal(serialisedValue, value); err != nil { - return nil, fmt.Errorf("failed to unmarshal value for request %s: %w", reqMetadata.RequestId, err) - } - - report = value.GetBytesValue() - if report == nil { - return nil, fmt.Errorf("failed to get report bytes for request %s", reqMetadata.RequestId) - } - } - - meta := ocrtypes.Metadata{ - Version: 1, - ExecutionID: reqMetadata.WorkflowExecutionId, - Timestamp: uint32(successOutcome.Timestamp.AsTime().Unix()), // nolint - DONID: reqMetadata.WorkflowDonId, - DONConfigVersion: reqMetadata.WorkflowDonConfigVersion, - WorkflowID: reqMetadata.WorkflowId, - WorkflowName: reqMetadata.WorkflowName, - WorkflowOwner: reqMetadata.WorkflowOwner, - ReportID: reqMetadata.ReportId, - } - - metadataPrepend, err := meta.Encode() - if err != nil { - return nil, fmt.Errorf("failed to encode metadata for request %s: %w", reqMetadata.RequestId, err) + if reqMetadata == nil { + // Without the metadata there is no request ID or key bundle to attribute a failure report to + r.lggr.Errorw("received successful consensus outcome without metadata, skipping", "seqNr", seqNr) + continue } - - reportWithMetaData := append(metadataPrepend, report...) - - // Check if the report is too large to transmit - if len(reportWithMetaData) > r.maxReportLengthBytes { - r.lggr.Errorw("report is too large to transmit", "seqNr", seqNr, "requestID", reqMetadata.RequestId, - "reportSize", len(reportWithMetaData), "maxReportLengthBytes", r.maxReportLengthBytes) - failureMsg := fmt.Sprintf( - "report too large: the report for this request is %d bytes which exceeds the maximum allowed size of %d bytes; reduce the size of the data being returned", - len(reportWithMetaData), r.maxReportLengthBytes) - info, err := createFailedConsensusReportInfo(reqMetadata.RequestId, reqMetadata.KeyBundleId, failureMsg, - oracletypes.ConsensusFailureCode_REPORT_TOO_LARGE) + r.lggr.Debugw("received successful consensus outcome", "seqNr", seqNr, "requestID", reqMetadata.RequestId) + + reportWithMetaData, failure := r.buildSuccessReport(successOutcome) + if failure != nil { + // The request is failed on its own so that the remaining outcomes in the round are still reported + r.lggr.Errorw("unable to build report for successful consensus outcome", "seqNr", seqNr, "requestID", reqMetadata.RequestId, + "failureCode", failure.code.String(), "failureMessage", failure.message) + info, err := createFailedConsensusReportInfo(reqMetadata.RequestId, reqMetadata.KeyBundleId, failure.message, failure.code) if err != nil { - return nil, fmt.Errorf("failed to create report info for oversized report %s: %w", reqMetadata.RequestId, err) + return nil, fmt.Errorf("failed to create report info for invalid consensus outcome %s: %w", reqMetadata.RequestId, err) } - reports = append(reports, ocr3types.ReportPlus[[]byte]{ - ReportWithInfo: ocr3types.ReportWithInfo[[]byte]{ - Report: []byte{}, - Info: info, - }, - TransmissionScheduleOverride: nil, - }) + + reports = append(reports, newFailureReport(info)) + failureIDs = append(failureIDs, reqMetadata.RequestId) continue } @@ -120,28 +90,114 @@ func (r *reportingPlugin) Reports(ctx context.Context, seqNr uint64, outcome ocr return nil, fmt.Errorf("failed to create report info for failed consensus outcome %s: %w", failedOutcome.RequestID, err) } - reports = append(reports, ocr3types.ReportPlus[[]byte]{ - ReportWithInfo: ocr3types.ReportWithInfo[[]byte]{ - Report: []byte{}, - Info: info, - }, - TransmissionScheduleOverride: nil, - }) + reports = append(reports, newFailureReport(info)) failureIDs = append(failureIDs, failedOutcome.RequestID) default: r.lggr.Warnw("received unknown consensus outcome type", "seqNr", seqNr, "outcome", outcome) } - - if len(reports) == r.maxNumberOfReports { - r.lggr.Warnw("maximum number of reports reached, stopping further report generation for this round", "seqNr", seqNr, "maxNumberOfReports", r.maxNumberOfReports) - break - } } r.lggr.Debugw("consensus plugin reports complete", "seqNr", seqNr, "numReports", len(reports), "successIDs", successIDs, "failureIDs", failureIDs) return reports, nil } +// reportFailure describes why a successful consensus outcome could not be turned into a report to sign. It is returned +// to the caller as a failed request in place of the report. +type reportFailure struct { + message string + code oracletypes.ConsensusFailureCode +} + +func newInvalidOutcomeFailure(format string, args ...any) *reportFailure { + return &reportFailure{ + message: fmt.Sprintf(format, args...), + code: oracletypes.ConsensusFailureCode_INVALID_OUTCOME, + } +} + +// buildSuccessReport builds the report to be signed for a successful consensus outcome, which is the encoded report +// metadata followed by the report payload. The outcome is validated first, as the DON attests to every byte of the +// report and a consumer cannot tell a report with an empty payload apart from a genuine result. +func (r *reportingPlugin) buildSuccessReport(successOutcome *oracletypes.ConsensusSuccessOutcome) ([]byte, *reportFailure) { + reqMetadata := successOutcome.Metadata + + var report []byte + switch reqMetadata.RequestType { + case oracletypes.RequestType_VALUE_CONSENSUS: + report = successOutcome.Outcome + case oracletypes.RequestType_REPORT_GENERATION: + // If the request type is report extract the report from the values.Value before signing it + value := &valuespb.Value{} + if err := proto.Unmarshal(successOutcome.Outcome, value); err != nil { + return nil, newInvalidOutcomeFailure("failed to unmarshal value for request %s: %v", reqMetadata.RequestId, err) + } + + report = value.GetBytesValue() + default: + return nil, newInvalidOutcomeFailure("unsupported request type %s for request %s", reqMetadata.RequestType, reqMetadata.RequestId) + } + + // A nil check alone would not catch a present but empty payload, such as a report request with an empty encoded payload + if len(report) == 0 { + return nil, newInvalidOutcomeFailure("consensus outcome for request %s has an empty report payload", reqMetadata.RequestId) + } + + if successOutcome.Timestamp == nil { + return nil, newInvalidOutcomeFailure("consensus outcome for request %s has no timestamp", reqMetadata.RequestId) + } + + // The report metadata carries the timestamp as a uint32 unix time, so a value outside of that range would otherwise + // be silently truncated. A zero timestamp is deliberately not rejected here: until every node sets + // include_error_observation_timestamps_flag the default value path can legitimately produce one, and all nodes must + // build identical reports for the same outcome. + unixTimestamp := successOutcome.Timestamp.AsTime().Unix() + if unixTimestamp < 0 || unixTimestamp > math.MaxUint32 { + return nil, newInvalidOutcomeFailure("consensus outcome for request %s has timestamp %d which does not fit in the report metadata", reqMetadata.RequestId, unixTimestamp) + } + + meta := ocrtypes.Metadata{ + Version: 1, + ExecutionID: reqMetadata.WorkflowExecutionId, + Timestamp: uint32(unixTimestamp), //nolint:gosec // G115 - range checked above + DONID: reqMetadata.WorkflowDonId, + DONConfigVersion: reqMetadata.WorkflowDonConfigVersion, + WorkflowID: reqMetadata.WorkflowId, + WorkflowName: reqMetadata.WorkflowName, + WorkflowOwner: reqMetadata.WorkflowOwner, + ReportID: reqMetadata.ReportId, + } + + metadataPrepend, err := meta.Encode() + if err != nil { + return nil, newInvalidOutcomeFailure("failed to encode metadata for request %s: %v", reqMetadata.RequestId, err) + } + + reportWithMetaData := append(metadataPrepend, report...) + + // Check if the report is too large to transmit + if len(reportWithMetaData) > r.maxReportLengthBytes { + return nil, &reportFailure{ + message: fmt.Sprintf( + "report too large: the report for this request is %d bytes which exceeds the maximum allowed size of %d bytes; reduce the size of the data being returned", + len(reportWithMetaData), r.maxReportLengthBytes), + code: oracletypes.ConsensusFailureCode_REPORT_TOO_LARGE, + } + } + + return reportWithMetaData, nil +} + +// newFailureReport creates a report with an empty body, as only the info is used to return a failure to the caller +func newFailureReport(info []byte) ocr3types.ReportPlus[[]byte] { + return ocr3types.ReportPlus[[]byte]{ + ReportWithInfo: ocr3types.ReportWithInfo[[]byte]{ + Report: []byte{}, + Info: info, + }, + TransmissionScheduleOverride: nil, + } +} + // The report info is created as a map else the OCR3OnchainKeyringMultiChainAdapter will not work. // OCR3OnchainKeyringMultiChainAdapter (in core) requires that the key bundle id is added to the map with the key // "keyBundleName". diff --git a/consensus/oracle/plugin/plugin_reports_test.go b/consensus/oracle/plugin/plugin_reports_test.go new file mode 100644 index 000000000..487c615cc --- /dev/null +++ b/consensus/oracle/plugin/plugin_reports_test.go @@ -0,0 +1,243 @@ +package plugin_test + +import ( + "strconv" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/require" + "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/types/known/structpb" + "google.golang.org/protobuf/types/known/timestamppb" + + ocrtypes "github.com/smartcontractkit/chainlink-common/pkg/capabilities/consensus/ocr3/types" + "github.com/smartcontractkit/chainlink-common/pkg/capabilities/consensus/requests" + "github.com/smartcontractkit/chainlink-common/pkg/logger" + + "github.com/smartcontractkit/chainlink-protos/cre/go/values" + + "github.com/smartcontractkit/libocr/offchainreporting2plus/ocr3types" + + "github.com/smartcontractkit/capabilities/consensus/metrics" + "github.com/smartcontractkit/capabilities/consensus/oracle" + "github.com/smartcontractkit/capabilities/consensus/oracle/plugin" + oracletypes "github.com/smartcontractkit/capabilities/consensus/oracle/types" +) + +// newReportsTestPlugin creates a reporting plugin for direct Reports() tests with the given report limits +func newReportsTestPlugin(t *testing.T, lggr logger.Logger, maxReportLengthBytes uint32, maxReportCount uint32) ocr3types.ReportingPlugin[[]byte] { + t.Helper() + + reqStore := requests.NewStore[*oracle.ConsensusRequest]() + metricsInstance, err := metrics.NewMetrics() + require.NoError(t, err) + + reportingPlugin, err := plugin.NewReportingPlugin(lggr, metricsInstance, f, n, reqStore, oracle.NewObservationQuorumTracker(), &ocrtypes.ReportingPluginConfig{ + MaxQueryLengthBytes: 1000000, + MaxObservationLengthBytes: 1000000, + MaxOutcomeLengthBytes: 1000000, + MaxReportLengthBytes: maxReportLengthBytes, + MaxReportCount: maxReportCount, + HistoricalOutcomeExpirySeqNrSpan: 5, + }, "evm", 1000) + require.NoError(t, err) + + return reportingPlugin +} + +func newReportsTestMetaData(requestID string, requestType oracletypes.RequestType) *oracletypes.RequestMetaData { + return &oracletypes.RequestMetaData{ + RequestId: requestID, + WorkflowExecutionId: "0102030405060708091011121314151617181920212223242526272829303132", + WorkflowId: "0039525c34de895c8fa68006bd63f6ce4a45ef1bc66377e791c6a8ae803dc0e4", + WorkflowOwner: "1139525c34de895c8fa68006bd634387a9f1192a", + WorkflowName: "a1b2c3d4e5f6a1b2c3d4", + WorkflowDonId: 1, + WorkflowDonConfigVersion: 1, + ReportId: "abcd", + KeyBundleId: "evm", + RequestType: requestType, + } +} + +func newSuccessOutcome(metadata *oracletypes.RequestMetaData, outcome []byte, timestamp *timestamppb.Timestamp) *oracletypes.ConsensusOutcome { + return &oracletypes.ConsensusOutcome{ + Outcome: &oracletypes.ConsensusOutcome_Success{ + Success: &oracletypes.ConsensusSuccessOutcome{ + Metadata: metadata, + Outcome: outcome, + Timestamp: timestamp, + }, + }, + } +} + +// serialiseValue serialises a value the way the outcome phase serialises a successful outcome +func serialiseValue(t *testing.T, value values.Value) []byte { + t.Helper() + serialisedValue, err := proto.MarshalOptions{Deterministic: true}.Marshal(values.Proto(value)) + require.NoError(t, err) + return serialisedValue +} + +func serialiseOutcome(t *testing.T, outcomes ...*oracletypes.ConsensusOutcome) []byte { + t.Helper() + serialisedOutcome, err := proto.MarshalOptions{Deterministic: true}.Marshal(&oracletypes.Outcome{Outcomes: outcomes}) + require.NoError(t, err) + return serialisedOutcome +} + +func reportInfoMap(t *testing.T, report ocr3types.ReportPlus[[]byte]) map[string]any { + t.Helper() + infos := &structpb.Struct{} + require.NoError(t, proto.Unmarshal(report.ReportWithInfo.Info, infos)) + return infos.AsMap() +} + +// verifyInvalidOutcomeReport checks that the report is a failure report with the INVALID_OUTCOME code for the request +func verifyInvalidOutcomeReport(t *testing.T, report ocr3types.ReportPlus[[]byte], expectedRequestID string, expectedFailureMessagePart string) { + t.Helper() + + require.Empty(t, report.ReportWithInfo.Report, "invalid outcome should have an empty report body") + + infoMap := reportInfoMap(t, report) + require.Equal(t, expectedRequestID, infoMap[plugin.InfoRequestID]) + require.Equal(t, "evm", infoMap[plugin.InfoKeyBundleName]) + require.Equal(t, oracletypes.ConsensusFailureCode_INVALID_OUTCOME.String(), infoMap[plugin.InfoConsensusFailureCode]) + require.Contains(t, infoMap[plugin.InfoConsensusFailureMessage], expectedFailureMessagePart) +} + +func Test_Reports_InvalidOutcome_ReturnsFailure(t *testing.T) { + lggr := logger.Test(t) + ctx := t.Context() + + reportingPlugin := newReportsTestPlugin(t, lggr, 10000, 100) + now := timestamppb.Now() + + testCases := []struct { + name string + outcome *oracletypes.ConsensusOutcome + expectedFailureMessagePart string + }{ + { + // This is what Report() puts on the wire for a ReportRequest with an empty EncodedPayload + name: "report generation with an empty payload", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType_REPORT_GENERATION), serialiseValue(t, values.NewBytes([]byte{})), now), + expectedFailureMessagePart: "empty report payload", + }, + { + name: "report generation with a value that is not bytes", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType_REPORT_GENERATION), serialiseValue(t, values.NewInt64(7)), now), + expectedFailureMessagePart: "empty report payload", + }, + { + name: "value consensus with an empty outcome", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType_VALUE_CONSENSUS), nil, now), + expectedFailureMessagePart: "empty report payload", + }, + { + name: "unknown request type", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType(99)), serialiseValue(t, values.NewBytes([]byte("payload"))), now), + expectedFailureMessagePart: "unsupported request type", + }, + { + name: "missing timestamp", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType_VALUE_CONSENSUS), serialiseValue(t, values.NewBytes([]byte("payload"))), nil), + expectedFailureMessagePart: "has no timestamp", + }, + { + name: "timestamp that does not fit in the report metadata", + outcome: newSuccessOutcome(newReportsTestMetaData("req", oracletypes.RequestType_VALUE_CONSENSUS), serialiseValue(t, values.NewBytes([]byte("payload"))), timestamppb.New(time.Unix(1<<33, 0))), + expectedFailureMessagePart: "does not fit in the report metadata", + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + reports, err := reportingPlugin.Reports(ctx, 1, serialiseOutcome(t, tc.outcome)) + require.NoError(t, err, "Reports should not return an error for an invalid outcome") + require.Len(t, reports, 1) + + verifyInvalidOutcomeReport(t, reports[0], "req", tc.expectedFailureMessagePart) + }) + } + + t.Run("metadata that cannot be encoded", func(t *testing.T) { + metadata := newReportsTestMetaData("req", oracletypes.RequestType_VALUE_CONSENSUS) + metadata.WorkflowExecutionId = "not-hex" + + reports, err := reportingPlugin.Reports(ctx, 1, serialiseOutcome(t, newSuccessOutcome(metadata, serialiseValue(t, values.NewBytes([]byte("payload"))), now))) + require.NoError(t, err) + require.Len(t, reports, 1) + + verifyInvalidOutcomeReport(t, reports[0], "req", "failed to encode metadata") + }) + + t.Run("missing metadata is skipped", func(t *testing.T) { + reports, err := reportingPlugin.Reports(ctx, 1, serialiseOutcome(t, newSuccessOutcome(nil, serialiseValue(t, values.NewBytes([]byte("payload"))), now))) + require.NoError(t, err) + require.Empty(t, reports, "an outcome without metadata cannot be attributed to a request") + }) +} + +func Test_Reports_InvalidOutcome_DoesNotFailRound(t *testing.T) { + lggr := logger.Test(t) + ctx := t.Context() + + reportingPlugin := newReportsTestPlugin(t, lggr, 10000, 100) + now := timestamppb.Now() + + // The first outcome cannot be reported, the second one is valid + invalidMetadata := newReportsTestMetaData("req-invalid", oracletypes.RequestType_VALUE_CONSENSUS) + invalidMetadata.WorkflowExecutionId = "not-hex" + validMetadata := newReportsTestMetaData("req-valid", oracletypes.RequestType_VALUE_CONSENSUS) + serialisedValue := serialiseValue(t, values.NewBytes([]byte("payload"))) + + reports, err := reportingPlugin.Reports(ctx, 1, serialiseOutcome(t, + newSuccessOutcome(invalidMetadata, serialisedValue, now), + newSuccessOutcome(validMetadata, serialisedValue, now), + )) + require.NoError(t, err, "Reports should not return an error when a single outcome is invalid") + require.Len(t, reports, 2, "Should have a report for each outcome") + + verifyInvalidOutcomeReport(t, reports[0], "req-invalid", "failed to encode metadata") + + // Verify the valid outcome is reported as a success with the expected metadata and payload + infoMap := reportInfoMap(t, reports[1]) + require.Equal(t, "req-valid", infoMap[plugin.InfoRequestID]) + require.Nil(t, infoMap[plugin.InfoConsensusFailureCode], "Should not have failure code for successful report") + + meta, payload, err := ocrtypes.Decode(reports[1].ReportWithInfo.Report) + require.NoError(t, err, "Failed to extract metadata fields from report") + require.Equal(t, uint32(now.AsTime().Unix()), meta.Timestamp) //nolint:gosec // G115 + require.Equal(t, serialisedValue, payload) + require.Len(t, reports[1].ReportWithInfo.Report, plugin.ReportMetaDataPrependLength+len(serialisedValue)) +} + +func Test_ReportCountLimit_IncludesOversizedReports(t *testing.T) { + lggr := logger.Test(t) + ctx := t.Context() + + // Create a reporting plugin with a small max report length and count + maxReportCount := uint32(3) + reportingPlugin := newReportsTestPlugin(t, lggr, 150, maxReportCount) + now := timestamppb.Now() + + // Create outcomes for more requests than the max report count, one of which produces an oversized report + var outcomes []*oracletypes.ConsensusOutcome + for i := 0; i < int(maxReportCount)+3; i++ { + data := "small" + if i == 1 { + data = strings.Repeat("x", 200) // More than 150 + } + outcomes = append(outcomes, newSuccessOutcome(newReportsTestMetaData("req-"+strconv.Itoa(i), oracletypes.RequestType_REPORT_GENERATION), + serialiseValue(t, values.NewBytes([]byte(data))), now)) + } + + // Call Reports and verify the oversized report counts towards the limit like any other report + reports, err := reportingPlugin.Reports(ctx, 1, serialiseOutcome(t, outcomes...)) + require.NoError(t, err) + require.Len(t, reports, int(maxReportCount)) + require.Equal(t, oracletypes.ConsensusFailureCode_REPORT_TOO_LARGE.String(), reportInfoMap(t, reports[1])[plugin.InfoConsensusFailureCode]) +} diff --git a/consensus/oracle/types/value_consensus_types.pb.go b/consensus/oracle/types/value_consensus_types.pb.go index 5a9192abd..0d9c4b34a 100644 --- a/consensus/oracle/types/value_consensus_types.pb.go +++ b/consensus/oracle/types/value_consensus_types.pb.go @@ -95,6 +95,9 @@ const ( // Indicates an error in the workflow logic, such that the workflow is not producing a consistent value type for the // same request across observers. ConsensusFailureCode_NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS ConsensusFailureCode = 7 + // The successful outcome could not be turned into a report to sign, e.g. it has an empty report payload, a missing or + // out of range timestamp, an unknown request type or metadata that cannot be encoded. Only the affected request fails. + ConsensusFailureCode_INVALID_OUTCOME ConsensusFailureCode = 8 ) // Enum value maps for ConsensusFailureCode. @@ -108,6 +111,7 @@ var ( 5: "MORE_THAN_ONE_VALID_OUTCOME_FOR_IDENTICAL_CONSENSUS", 6: "NO_VALUES_MET_FPLUS1_THRESHOLD_FOR_IDENTICAL_CONSENSUS", 7: "NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS", + 8: "INVALID_OUTCOME", } ConsensusFailureCode_value = map[string]int32{ "CONSENSUS_CALCULATION_FAILED": 0, @@ -118,6 +122,7 @@ var ( "MORE_THAN_ONE_VALID_OUTCOME_FOR_IDENTICAL_CONSENSUS": 5, "NO_VALUES_MET_FPLUS1_THRESHOLD_FOR_IDENTICAL_CONSENSUS": 6, "NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS": 7, + "INVALID_OUTCOME": 8, } ) @@ -323,6 +328,7 @@ type RequestObservation struct { ReceivedAt *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=received_at,json=receivedAt,proto3" json:"received_at,omitempty"` RemoveLibUseInFailureMessageFormattingFlag bool `protobuf:"varint,4,opt,name=remove_lib_use_in_failure_message_formatting_flag,json=removeLibUseInFailureMessageFormattingFlag,proto3" json:"remove_lib_use_in_failure_message_formatting_flag,omitempty"` // remove use of libraries in failure message formatting; flag to be removed after rollout UpdateErrorHandlingFlag bool `protobuf:"varint,5,opt,name=update_error_handling_flag,json=updateErrorHandlingFlag,proto3" json:"update_error_handling_flag,omitempty"` // migrate system errors to user errors; flag to be removed after rollout + IncludeErrorObservationTimestampsFlag bool `protobuf:"varint,6,opt,name=include_error_observation_timestamps_flag,json=includeErrorObservationTimestampsFlag,proto3" json:"include_error_observation_timestamps_flag,omitempty"` // include error observations in the median outcome timestamp; flag to be removed after rollout unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -392,6 +398,13 @@ func (x *RequestObservation) GetUpdateErrorHandlingFlag() bool { return false } +func (x *RequestObservation) GetIncludeErrorObservationTimestampsFlag() bool { + if x != nil { + return x.IncludeErrorObservationTimestampsFlag + } + return false +} + type Observation struct { state protoimpl.MessageState `protogen:"open.v1"` Observations map[string]*RequestObservation `protobuf:"bytes,1,rep,name=observations,proto3" json:"observations,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` @@ -828,14 +841,15 @@ const file_value_consensus_types_proto_rawDesc = "" + "\x05Query\x12\x1e\n" + "\n" + "requestIDs\x18\x01 \x03(\tR\n" + - "requestIDs\"\xf3\x02\n" + + "requestIDs\"\xcd\x03\n" + "\x12RequestObservation\x12B\n" + "\bmetadata\x18\x01 \x01(\v2&.value_consensus_types.RequestMetaDataR\bmetadata\x128\n" + "\x05input\x18\x02 \x01(\v2\".sdk.v1alpha.SimpleConsensusInputsR\x05input\x12;\n" + "\vreceived_at\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\n" + "receivedAt\x12e\n" + "1remove_lib_use_in_failure_message_formatting_flag\x18\x04 \x01(\bR*removeLibUseInFailureMessageFormattingFlag\x12;\n" + - "\x1aupdate_error_handling_flag\x18\x05 \x01(\bR\x17updateErrorHandlingFlag\"\xd3\x01\n" + + "\x1aupdate_error_handling_flag\x18\x05 \x01(\bR\x17updateErrorHandlingFlag\x12X\n" + + ")include_error_observation_timestamps_flag\x18\x06 \x01(\bR%includeErrorObservationTimestampsFlag\"\xd3\x01\n" + "\vObservation\x12X\n" + "\fobservations\x18\x01 \x03(\v24.value_consensus_types.Observation.ObservationsEntryR\fobservations\x1aj\n" + "\x11ObservationsEntry\x12\x10\n" + @@ -868,7 +882,7 @@ const file_value_consensus_types_proto_rawDesc = "" + "\x05value\x18\x02 \x01(\x04R\x05value*9\n" + "\vRequestType\x12\x13\n" + "\x0fVALUE_CONSENSUS\x10\x00\x12\x15\n" + - "\x11REPORT_GENERATION\x10\x01*\xda\x02\n" + + "\x11REPORT_GENERATION\x10\x01*\xef\x02\n" + "\x14ConsensusFailureCode\x12 \n" + "\x1cCONSENSUS_CALCULATION_FAILED\x10\x00\x12%\n" + "!FAILED_TO_CALCULATE_CONSENSUS_MDD\x10\x01\x12\x1a\n" + @@ -877,7 +891,8 @@ const file_value_consensus_types_proto_rawDesc = "" + "\x10REPORT_TOO_LARGE\x10\x04\x127\n" + "3MORE_THAN_ONE_VALID_OUTCOME_FOR_IDENTICAL_CONSENSUS\x10\x05\x12:\n" + "6NO_VALUES_MET_FPLUS1_THRESHOLD_FOR_IDENTICAL_CONSENSUS\x10\x06\x12;\n" + - "7NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS\x10\aB\x18Z\x16consensus/oracle/typesb\x06proto3" + "7NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS\x10\a\x12\x13\n" + + "\x0fINVALID_OUTCOME\x10\bB\x18Z\x16consensus/oracle/typesb\x06proto3" var ( file_value_consensus_types_proto_rawDescOnce sync.Once diff --git a/consensus/oracle/types/value_consensus_types.proto b/consensus/oracle/types/value_consensus_types.proto index e8654a9bf..26cd45719 100644 --- a/consensus/oracle/types/value_consensus_types.proto +++ b/consensus/oracle/types/value_consensus_types.proto @@ -40,6 +40,7 @@ message RequestObservation { bool remove_lib_use_in_failure_message_formatting_flag = 4; // remove use of libraries in failure message formatting; flag to be removed after rollout bool update_error_handling_flag = 5; // migrate system errors to user errors; flag to be removed after rollout + bool include_error_observation_timestamps_flag = 6; // include error observations in the median outcome timestamp; flag to be removed after rollout } message Observation { @@ -89,6 +90,9 @@ enum ConsensusFailureCode { // Indicates an error in the workflow logic, such that the workflow is not producing a consistent value type for the // same request across observers. NO_SINGLE_VALUE_TYPE_MET_FPLUS1_THRESHOLD_FOR_CONSENSUS = 7; + // The successful outcome could not be turned into a report to sign, e.g. it has an empty report payload, a missing or + // out of range timestamp, an unknown request type or metadata that cannot be encoded. Only the affected request fails. + INVALID_OUTCOME = 8; } message ConsensusFailedOutcome {