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
3 changes: 2 additions & 1 deletion chain_capabilities/evm/actions/actions.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ import (
"github.com/smartcontractkit/capabilities/chain_capabilities/evm/internal/contracts"
"github.com/smartcontractkit/capabilities/chain_capabilities/evm/metering"
"github.com/smartcontractkit/capabilities/chain_capabilities/evm/monitoring"
commonMon "github.com/smartcontractkit/capabilities/libs/monitoring"
)

type ConsensusHandler interface {
Expand Down Expand Up @@ -100,7 +101,7 @@ func (e *EVM) initLimiters(limitsFactory limits.Factory) (err error) {
}

func requestID(meta capabilities.RequestMetadata) string {
return meta.WorkflowExecutionID + ":" + meta.ReferenceID
return commonMon.RequestID(meta.WorkflowExecutionID, meta.ReferenceID)
}

func (e *EVM) CallContract(
Expand Down
43 changes: 39 additions & 4 deletions chain_capabilities/evm/actions/write_report.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,13 @@ import (
"github.com/ethereum/go-ethereum/common"
"github.com/jpillora/backoff"

"github.com/smartcontractkit/chainlink-common/pkg/beholder"
"github.com/smartcontractkit/chainlink-common/pkg/capabilities"
ocrtypes "github.com/smartcontractkit/chainlink-common/pkg/capabilities/consensus/ocr3/types"
"github.com/smartcontractkit/chainlink-common/pkg/capabilities/v2/chain-capabilities/evm"
commoncfg "github.com/smartcontractkit/chainlink-common/pkg/config"
"github.com/smartcontractkit/chainlink-common/pkg/settings/limits"
"github.com/smartcontractkit/chainlink-common/pkg/types"
"github.com/smartcontractkit/chainlink-common/pkg/utils/retry"

"github.com/smartcontractkit/chainlink-common/pkg/logger"
Expand All @@ -42,6 +45,20 @@ func decodeReportMetadata(data []byte) (ocrtypes.Metadata, error) {
return metadata, err
}

type WriteReport struct {
types.EVMService
forwarderClient contracts.CREForwarderClient
ReceiverGasMinimum uint64
chainSelector uint64

lggr logger.Logger
beholderProcessor beholder.ProtoProcessor
messageBuilder *monitoring.MessageBuilder

txGasLimit limits.BoundLimiter[uint64]
reportSizeLimit limits.BoundLimiter[commoncfg.Size]
}

func (e *EVM) WriteReport(ctx context.Context, metadata capabilities.RequestMetadata, input *evm.WriteReportRequest) (*capabilities.ResponseAndMetadata[*evm.WriteReportReply], error) {
ctx = metadata.ContextWithCRE(ctx)
telemetryContext := monitoring.TelemetryContext{TsStart: time.Now().UnixMilli(), RequestMetadata: metadata}
Expand All @@ -66,7 +83,25 @@ func (e *EVM) WriteReport(ctx context.Context, metadata capabilities.RequestMeta
return &responseAndMetadata, nil
}

func (e *EVM) getFee(ctx context.Context, txIdempotencyKey string) (*big.Float, error) {
func (e *EVM) executeWriteReport(ctx context.Context, request *evm.WriteReportRequest, metadata capabilities.RequestMetadata, telemetryContext monitoring.TelemetryContext) (*evm.WriteReportReply, capabilities.ResponseMetadata, error) {
wr := &WriteReport{
EVMService: e.EVMService,
forwarderClient: e.forwarderClient,
ReceiverGasMinimum: e.ReceiverGasMinimum,
chainSelector: e.chainSelector,

lggr: e.messageBuilder.RequestLggr(e.lggr, telemetryContext),
beholderProcessor: e.beholderProcessor,
messageBuilder: e.messageBuilder,

txGasLimit: e.txGasLimit,
reportSizeLimit: e.reportSizeLimit,
}

return wr.executeWriteReport(ctx, request, metadata, telemetryContext)
}

func (e *WriteReport) getFee(ctx context.Context, txIdempotencyKey string) (*big.Float, error) {
if txIdempotencyKey == "" {
return nil, fmt.Errorf("txIdempotencyKey is empty, cannot retrieve transaction fee")
}
Expand All @@ -80,7 +115,7 @@ func (e *EVM) getFee(ctx context.Context, txIdempotencyKey string) (*big.Float,
return feeInEth, nil
}

func (e *EVM) executeWriteReport(ctx context.Context, request *evm.WriteReportRequest, metadata capabilities.RequestMetadata, telemetryContext monitoring.TelemetryContext) (*evm.WriteReportReply, capabilities.ResponseMetadata, error) {
func (e *WriteReport) executeWriteReport(ctx context.Context, request *evm.WriteReportRequest, metadata capabilities.RequestMetadata, telemetryContext monitoring.TelemetryContext) (*evm.WriteReportReply, capabilities.ResponseMetadata, error) {
transmissionID, err := getTransmissionID(metadata.WorkflowExecutionID, request)
if err != nil {
return nil, capabilities.ResponseMetadata{}, err
Expand Down Expand Up @@ -216,7 +251,7 @@ func getInvalidStateErrorMessage(state uint8) string {
return fmt.Sprintf("unexpected transmission state: %v", state)
}

func (e *EVM) processUnrecoverableTxState(ctx context.Context, request *evm.WriteReportRequest, metadata capabilities.RequestMetadata, txHash evmtypes.Hash, transmissionInfo contracts.TransmissionInfo, transmissionID contracts.TransmissionID, txAttemptedLocally bool) (*evm.WriteReportReply, error) {
func (e *WriteReport) processUnrecoverableTxState(ctx context.Context, request *evm.WriteReportRequest, metadata capabilities.RequestMetadata, txHash evmtypes.Hash, transmissionInfo contracts.TransmissionInfo, transmissionID contracts.TransmissionID, txAttemptedLocally bool) (*evm.WriteReportReply, error) {
if !txAttemptedLocally {
e.lggr.Infow("returning without a transmission attempt - transmission already attempted, receiver was marked as invalid", "executionID", metadata.WorkflowExecutionID)
} else {
Expand Down Expand Up @@ -259,7 +294,7 @@ func getTransmissionID(workflowExecutionID string, request *evm.WriteReportReque
return transmissionID, nil
}

func (e *EVM) fetchTransactionReceiptAndCreateReply(ctx context.Context, txHash evmtypes.Hash, receiverStatus evm.ReceiverContractExecutionStatus, errorMessage *string) (*evm.WriteReportReply, error) {
func (e *WriteReport) fetchTransactionReceiptAndCreateReply(ctx context.Context, txHash evmtypes.Hash, receiverStatus evm.ReceiverContractExecutionStatus, errorMessage *string) (*evm.WriteReportReply, error) {
// TODO: PLEX-1524 - we need retry logic here in case the underlying RPC is lagging behind the one that submitted the TX.
txReceipt, err := e.EVMService.GetTransactionReceipt(ctx, evmtypes.GeTransactionReceiptRequest{
Hash: txHash,
Expand Down
2 changes: 1 addition & 1 deletion chain_capabilities/evm/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ go 1.25.3
require (
github.com/ethereum/go-ethereum v1.16.2
github.com/google/go-cmp v0.7.0
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16
github.com/smartcontractkit/chain-selectors v1.0.67
github.com/smartcontractkit/chainlink-common v0.9.6-0.20251028165938-6226fd06e3f8
github.com/smartcontractkit/chainlink-evm v0.3.4-0.20251020152820-5fb041bf92b7
Expand Down
4 changes: 2 additions & 2 deletions chain_capabilities/evm/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -457,8 +457,8 @@ github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp
github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME=
github.com/sirupsen/logrus v1.4.1/go.mod h1:ni0Sbl8bgC9z8RoU9G6nDWqqs/fq4eDPysMBDgk/93Q=
github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91 h1:Cm4uo5A9vIpovV99Dm2aMcVRdmVbinnwp43v4DEcYhE=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91/go.mod h1:CmCb4hm2aAzdQqwZFGBc2/37JEpqDBix4CWr2OO7ls0=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16 h1:9vrw/5SFeEvIMmqnzlsiLCfJYU14PK1b7rmMyS4GFlE=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16/go.mod h1:CmCb4hm2aAzdQqwZFGBc2/37JEpqDBix4CWr2OO7ls0=
github.com/smartcontractkit/chain-selectors v1.0.67 h1:gxTqP/JC40KDe3DE1SIsIKSTKTZEPyEU1YufO1admnw=
github.com/smartcontractkit/chain-selectors v1.0.67/go.mod h1:xsKM0aN3YGcQKTPRPDDtPx2l4mlTN1Djmg0VVXV40b8=
github.com/smartcontractkit/chainlink-common v0.9.6-0.20251028165938-6226fd06e3f8 h1:CIIY+VF5ZDnodC2cdviGZE4Ze6kvaaxtqZbVI7x+32w=
Expand Down
7 changes: 7 additions & 0 deletions chain_capabilities/evm/monitoring/messages.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"google.golang.org/protobuf/proto"

evmcap "github.com/smartcontractkit/chainlink-common/pkg/capabilities/v2/chain-capabilities/evm"
"github.com/smartcontractkit/chainlink-common/pkg/logger"
sdkpb "github.com/smartcontractkit/chainlink-protos/cre/go/sdk"
valuespb "github.com/smartcontractkit/chainlink-protos/cre/go/values/pb"

Expand Down Expand Up @@ -54,6 +55,12 @@ func NewMessageBuilder(chainInfo types.ChainInfo, capInfo capabilities.Capabilit
return &MessageBuilder{ChainInfo: chainInfo, CapInfo: capInfo, nodeAddress: nodeAddress}
}

func (m *MessageBuilder) RequestLggr(lggr logger.SugaredLogger, telemetryContext TelemetryContext) logger.SugaredLogger {
attrs := m.BuildExecutionContext(telemetryContext).LogAttributes()
lggrAttrs := attrsToErrorKV(attrs)
return lggr.With(lggrAttrs...)
}

func (m *MessageBuilder) BuildCallContractInitiated(tc TelemetryContext, msg *evm.CallMsg, bn int64) *CallContractInitiated {
return &CallContractInitiated{Req: &CallContractRequest{BlockNumber: bn, ContractAddress: common.Bytes2Hex(msg.To[:])}, ExecutionContext: m.BuildExecutionContext(tc)}
}
Expand Down
2 changes: 1 addition & 1 deletion integration_tests/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -296,7 +296,7 @@ require (
github.com/shirou/gopsutil/v3 v3.24.3 // indirect
github.com/shopspring/decimal v1.4.0 // indirect
github.com/sigurn/crc16 v0.0.0-20211026045750-20ab5afb07e3 // indirect
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91 // indirect
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16 // indirect
github.com/smartcontractkit/chainlink-aptos v0.0.0-20251027153600-2b072ff3618e // indirect
github.com/smartcontractkit/chainlink-automation v0.8.1 // indirect
github.com/smartcontractkit/chainlink-ccip v0.1.1-solana.0.20251024071356-520275eaaf00 // indirect
Expand Down
4 changes: 2 additions & 2 deletions integration_tests/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -1083,8 +1083,8 @@ github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrf
github.com/sirupsen/logrus v1.8.1/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0=
github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ=
github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91 h1:Cm4uo5A9vIpovV99Dm2aMcVRdmVbinnwp43v4DEcYhE=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251023190209-61f60e422f91/go.mod h1:CmCb4hm2aAzdQqwZFGBc2/37JEpqDBix4CWr2OO7ls0=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16 h1:9vrw/5SFeEvIMmqnzlsiLCfJYU14PK1b7rmMyS4GFlE=
github.com/smartcontractkit/capabilities/libs v0.0.0-20251029151227-b0136f26ab16/go.mod h1:CmCb4hm2aAzdQqwZFGBc2/37JEpqDBix4CWr2OO7ls0=
github.com/smartcontractkit/chain-selectors v1.0.75 h1:72csyj5UL0Agi81gIX6QWGfGrRmUm3dSh/2nLCpUr+g=
github.com/smartcontractkit/chain-selectors v1.0.75/go.mod h1:xsKM0aN3YGcQKTPRPDDtPx2l4mlTN1Djmg0VVXV40b8=
github.com/smartcontractkit/chainlink-aptos v0.0.0-20251027153600-2b072ff3618e h1:HIgcJV/CyhBntE5gK/8WitVzqD0k8PkuYj+lhfa6B6U=
Expand Down
4 changes: 4 additions & 0 deletions libs/monitoring/common.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,3 +94,7 @@ func (m *MetricsCapBasic) RecordEmit(ctx context.Context, start, emit uint64, at
m.capTimestampEmit.Record(ctx, int64(emit), attrs)
m.capDuration.Record(ctx, int64(emit-start), attrs)
}

func RequestID(workflowExecutionID, reference string) string {
return workflowExecutionID + ":" + reference
}
6 changes: 2 additions & 4 deletions libs/monitoring/metadata.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,13 +58,11 @@ func (x *ExecutionContext) LogAttributes() []attribute.KeyValue {
// Execution Context - Workflow (capabilities.RequestMetadata)
attribute.String("workflow_id", ValOrUnknown(x.GetMetaWorkflowId())),
attribute.String("workflow_owner", ValOrUnknown(x.GetMetaWorkflowOwner())),
// Notice: We lower the cardinality on the WorkflowExecutionID so it can be used by metrics
// This label has good chances to be unique per workflow, in a reasonable bounded time window
// TODO: enable this when sufficiently tested (PromQL queries like alerts might need to change if this is used)
//attribute.String("workflow_execution_id_short", ValShortOrUnknown(x.GetMetaWorkflowExecutionId(), WorkflowExecutionIDShortLen)),
attribute.String("workflow_execution_id", ValOrUnknown(x.GetMetaWorkflowExecutionId())),
attribute.String("workflow_name", ValOrUnknown(workflowName)),
attribute.Int64("workflow_don_config_version", int64(x.GetMetaWorkflowDonConfigVersion())),
attribute.String("reference_id", ValOrUnknown(x.GetMetaReferenceId())),
attribute.String("request_id", RequestID(ValOrUnknown(x.GetMetaWorkflowExecutionId()), ValOrUnknown(x.GetMetaReferenceId()))),
)
}

Expand Down
Loading