Skip to content
Draft
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
18 changes: 18 additions & 0 deletions consensus/action/capability.go
Original file line number Diff line number Diff line change
Expand Up @@ -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:]

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}
9 changes: 9 additions & 0 deletions consensus/action/capability_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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)))
}
7 changes: 7 additions & 0 deletions consensus/oracle/plugin/errors_consensus_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}},
}

Expand Down
1 change: 1 addition & 0 deletions consensus/oracle/plugin/plugin_observation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
13 changes: 13 additions & 0 deletions consensus/oracle/plugin/plugin_outcome.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,7 @@
}

// addRequestOutcomeToBatch adds the outcome for a single request to the outcome batch. Returns false if batch does not have capacity to add the outcome.
func (r *reportingPlugin) addRequestOutcomeToBatch(ctx context.Context, lggr logger.Logger, requestID string, observations []*oracletypes.RequestObservation, outcome *batching.OutcomeBatch) (bool, error) {

Check warning on line 98 in consensus/oracle/plugin/plugin_outcome.go

View check run for this annotation

CL-sonarqube-production / SonarQube Code Analysis

Refactor this method to reduce its Cognitive Complexity from 36 to the 30 allowed.

[S3776] Cognitive Complexity of functions should not be too high See more on https://sonarqube.main.prod.cldev.sh/project/issues?id=smartcontractkit_capabilities&pullRequest=748&issues=3aee9722-2d2e-424c-8f87-82e8f9c000bc&open=3aee9722-2d2e-424c-8f87-82e8f9c000bc
// false is ok to use as the default for the updateErrorHandlingFlag parameter as the flag pertains to how the error is reported when observations have different types,
// in this case we know that all the observations will be of type []byte so the error will not occur and thus the flag will not have an effect on the outcome.
consensusMDD, err := r.calculateConsensusMetadataDescriptorAndDefault(lggr, observations, false)
Expand All @@ -108,9 +108,11 @@
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
Expand All @@ -120,6 +122,10 @@
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)
Expand All @@ -145,9 +151,16 @@
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 := &timestamppb.Timestamp{}
if len(timestamps) > 0 {
timestamp = calculateMedianTimestamp(timestamps)
Expand Down
112 changes: 112 additions & 0 deletions consensus/oracle/plugin/plugin_outcome_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{}
Expand Down Expand Up @@ -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) {
Expand Down
Loading
Loading