diff --git a/chain_capabilities/evm/actions/actions.go b/chain_capabilities/evm/actions/actions.go index 857fc1bd3..34dca0bc3 100644 --- a/chain_capabilities/evm/actions/actions.go +++ b/chain_capabilities/evm/actions/actions.go @@ -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 { @@ -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( diff --git a/chain_capabilities/evm/actions/write_report.go b/chain_capabilities/evm/actions/write_report.go index 4c849559d..37eb50839 100644 --- a/chain_capabilities/evm/actions/write_report.go +++ b/chain_capabilities/evm/actions/write_report.go @@ -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" @@ -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} @@ -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") } @@ -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 @@ -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 { @@ -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, diff --git a/chain_capabilities/evm/go.mod b/chain_capabilities/evm/go.mod index 8db927858..f0adea876 100644 --- a/chain_capabilities/evm/go.mod +++ b/chain_capabilities/evm/go.mod @@ -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 diff --git a/chain_capabilities/evm/go.sum b/chain_capabilities/evm/go.sum index cb1e7aa46..07b8b6047 100644 --- a/chain_capabilities/evm/go.sum +++ b/chain_capabilities/evm/go.sum @@ -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= diff --git a/chain_capabilities/evm/monitoring/messages.go b/chain_capabilities/evm/monitoring/messages.go index f37d79fe3..77b4967ab 100644 --- a/chain_capabilities/evm/monitoring/messages.go +++ b/chain_capabilities/evm/monitoring/messages.go @@ -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" @@ -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)} } diff --git a/integration_tests/go.mod b/integration_tests/go.mod index d3c52911c..cb2e191de 100644 --- a/integration_tests/go.mod +++ b/integration_tests/go.mod @@ -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 diff --git a/integration_tests/go.sum b/integration_tests/go.sum index e5e2c82ce..73c19f55c 100644 --- a/integration_tests/go.sum +++ b/integration_tests/go.sum @@ -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= diff --git a/libs/monitoring/common.go b/libs/monitoring/common.go index 4fce33621..a76d29ab1 100644 --- a/libs/monitoring/common.go +++ b/libs/monitoring/common.go @@ -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 +} diff --git a/libs/monitoring/metadata.go b/libs/monitoring/metadata.go index 0ea76330f..18da1066e 100644 --- a/libs/monitoring/metadata.go +++ b/libs/monitoring/metadata.go @@ -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()))), ) }