diff --git a/chain_capabilities/common/gasmeter/gasmeter.go b/chain_capabilities/common/gasmeter/gasmeter.go new file mode 100644 index 000000000..ce5953461 --- /dev/null +++ b/chain_capabilities/common/gasmeter/gasmeter.go @@ -0,0 +1,100 @@ +// Package gasmeter emits the cre:workflow:gas: usage MeterRecord +// that bills one chain write. Every chain-write capability (EVM, Solana, ...) +// uses it so the record shape, identity and the reconciler log line stay +// identical across chains. +package gasmeter + +import ( + "context" + "math/big" + "strconv" + + "github.com/smartcontractkit/chainlink-common/pkg/capabilities" + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/resourcemanager" +) + +// Meter emits one METER_ACTION_USAGE record per chain write. A nil *Meter is a +// no-op everywhere, so callers only need nil checks when registering it as a +// service. +type Meter struct { + lggr logger.Logger + rm *resourcemanager.ResourceManager + identity resourcemanager.ResourceIdentity + resourceType string +} + +// New builds a Meter from the LOOP metering config ([Metering] on the host). +// It returns nil when MeterRecordsEnabled is false: gas usage records share +// that gate with durable resource metering. Records are never snapshotted. +// capabilityDonID of 0 means unknown and leaves the DON id unset. +func New(lggr logger.Logger, cfg resourcemanager.Config, chainSelector uint64, capabilityDonID uint32) *Meter { + if !cfg.MeterRecordsEnabled { + return nil + } + rmCfg := cfg.ResourceManagerConfig + rmCfg.MeterSnapshotsEnabled = false + resourceType := resourcemanager.WorkflowGasResourceType(chainSelector) + identity := resourcemanager.WithWorkflowUsagePool( + resourcemanager.NewBaseIdentity(cfg.DeploymentIdentity, resourcemanager.EmittingServiceChainWrite, ""), + resourceType, + ) + if capabilityDonID != 0 { + identity = identity.WithDonID(strconv.FormatUint(uint64(capabilityDonID), 10)) + } + if rmCfg.Emitter == nil { + lggr.Errorw("Capability usage metering enabled but this LOOP has no durable emitter; gas usage records will not be delivered") + } + return &Meter{ + lggr: lggr, + rm: resourcemanager.NewResourceManager(lggr, rmCfg), + identity: identity, + resourceType: resourceType, + } +} + +// Start starts the underlying ResourceManager. +func (m *Meter) Start(ctx context.Context) error { + if m == nil { + return nil + } + return m.rm.Start(ctx) +} + +// Close closes the underlying ResourceManager. +func (m *Meter) Close() error { + if m == nil { + return nil + } + return m.rm.Close() +} + +// Emit emits the gas usage record for one chain write and logs the emission. +// txHash (or tx signature) is the capability event id: one record per +// on-chain write, identical on every node of the DON that observes it. The +// log line is a contract consumed by the billing reconciler (fields: +// executionID, eventID, resourceType, value, orgID, txHash) and must stay +// stable. Fail-open: never affects the reply. +func (m *Meter) Emit(ctx context.Context, metadata capabilities.RequestMetadata, txHash string, fee *big.Int) { + if m == nil || fee == nil { + return + } + resourceID, err := resourcemanager.WorkflowUsageResourceID(metadata.WorkflowID, metadata.WorkflowExecutionID) + if err != nil { + m.lggr.Errorw("Gas usage meter record not emitted", "err", err, "executionID", metadata.WorkflowExecutionID) + return + } + m.rm.EmitUsageValue(ctx, m.identity, txHash, fee, resourcemanager.UtilizationFields{ + ResourceType: m.resourceType, + ResourceID: resourceID, + OrgID: metadata.OrgID, + }) + m.lggr.Infow("Emitted capability usage meter record", + "executionID", metadata.WorkflowExecutionID, + "eventID", txHash, + "resourceType", m.resourceType, + "value", fee.String(), + "orgID", metadata.OrgID, + "txHash", txHash, + ) +} diff --git a/chain_capabilities/common/gasmeter/gasmeter_test.go b/chain_capabilities/common/gasmeter/gasmeter_test.go new file mode 100644 index 000000000..832f5ea0f --- /dev/null +++ b/chain_capabilities/common/gasmeter/gasmeter_test.go @@ -0,0 +1,118 @@ +package gasmeter + +import ( + "context" + "math/big" + "testing" + + "github.com/stretchr/testify/require" + "go.uber.org/zap/zapcore" + "google.golang.org/protobuf/proto" + + "github.com/smartcontractkit/chainlink-common/pkg/capabilities" + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/resourcemanager" + meteringpb "github.com/smartcontractkit/chainlink-protos/metering/go" +) + +type recordingEmitter struct{ records []*meteringpb.MeterRecord } + +func (r *recordingEmitter) Emit(_ context.Context, body []byte, _ ...any) error { + var rec meteringpb.MeterRecord + if err := proto.Unmarshal(body, &rec); err != nil { + return err + } + r.records = append(r.records, &rec) + return nil +} + +func enabledCfg(emitter resourcemanager.Emitter) resourcemanager.Config { + return resourcemanager.Config{ + MeterRecordsEnabled: true, + MeterSnapshotsEnabled: true, // must be forced off for usage records + Emitter: emitter, + DeploymentIdentity: resourcemanager.DeploymentIdentity{Product: "cre", Tenant: "t", Environment: "test", Zone: "z", NodeID: "node-1"}, + } +} + +func TestNew(t *testing.T) { + t.Run("nil when MeterRecordsEnabled is false", func(t *testing.T) { + require.Nil(t, New(logger.Test(t), resourcemanager.Config{Emitter: &recordingEmitter{}}, 1, 7)) + var m *Meter + require.NoError(t, m.Start(t.Context())) + require.NoError(t, m.Close()) + m.Emit(t.Context(), capabilities.RequestMetadata{}, "0x", big.NewInt(1)) // no panic + }) + + t.Run("logs an error when enabled without an emitter", func(t *testing.T) { + lggr, obs := logger.TestObserved(t, zapcore.ErrorLevel) + cfg := enabledCfg(nil) + require.NotNil(t, New(lggr, cfg, 1, 7)) + require.Len(t, obs.FilterMessage("Capability usage metering enabled but this LOOP has no durable emitter; gas usage records will not be delivered").All(), 1) + }) +} + +func TestEmit(t *testing.T) { + metadata := capabilities.RequestMetadata{WorkflowID: "wf-1", WorkflowExecutionID: "exec-1", OrgID: "org-1"} + fee, _ := new(big.Int).SetString("123456789012345678901234", 10) + + t.Run("emits one gas record keyed by tx hash and logs the contract line", func(t *testing.T) { + lggr, obs := logger.TestObserved(t, zapcore.InfoLevel) + emitter := &recordingEmitter{} + m := New(lggr, enabledCfg(emitter), 421614, 7) + require.NotNil(t, m) + + m.Emit(t.Context(), metadata, "0xabc", fee) + + require.Len(t, emitter.records, 1) + rec := emitter.records[0] + require.Equal(t, meteringpb.MeterAction_METER_ACTION_USAGE, rec.GetAction()) + id := rec.GetIdentity() + require.Equal(t, resourcemanager.EmittingServiceChainWrite, id.GetService()) + require.Equal(t, "cre:workflow:gas", id.GetResourcePool()) + require.Equal(t, "cre:workflow:gas:421614", id.GetResourcePoolId()) + require.Equal(t, "7", id.GetDon().GetDonId()) + require.Equal(t, "node-1", id.GetDon().GetNodeId()) + require.Len(t, rec.GetUtilizations(), 1) + u := rec.GetUtilizations()[0] + require.Equal(t, "cre:workflow:gas:421614", u.GetResourceType()) + require.Equal(t, "wf-1:exec-1", u.GetResourceId()) + require.Equal(t, "0xabc", u.GetEventId()) + require.Equal(t, "org-1", u.GetOrgId()) + require.Equal(t, "123456789012345678901234", u.GetValue()) + + logs := obs.FilterMessage("Emitted capability usage meter record").All() + require.Len(t, logs, 1) + fields := logs[0].ContextMap() + require.Equal(t, "exec-1", fields["executionID"]) + require.Equal(t, "0xabc", fields["eventID"]) + require.Equal(t, "cre:workflow:gas:421614", fields["resourceType"]) + require.Equal(t, "123456789012345678901234", fields["value"]) + require.Equal(t, "org-1", fields["orgID"]) + require.Equal(t, "0xabc", fields["txHash"]) + }) + + t.Run("unknown DON id leaves the identity DON unset", func(t *testing.T) { + emitter := &recordingEmitter{} + m := New(logger.Test(t), enabledCfg(emitter), 1, 0) + m.Emit(t.Context(), metadata, "0xabc", fee) + require.Len(t, emitter.records, 1) + require.Empty(t, emitter.records[0].GetIdentity().GetDon().GetDonId()) + }) + + t.Run("no-op without a fee", func(t *testing.T) { + emitter := &recordingEmitter{} + m := New(logger.Test(t), enabledCfg(emitter), 1, 7) + m.Emit(t.Context(), metadata, "0xabc", nil) + require.Empty(t, emitter.records) + }) + + t.Run("refuses a malformed resource id instead of emitting a bad record", func(t *testing.T) { + lggr, obs := logger.TestObserved(t, zapcore.ErrorLevel) + emitter := &recordingEmitter{} + m := New(lggr, enabledCfg(emitter), 1, 7) + m.Emit(t.Context(), capabilities.RequestMetadata{WorkflowID: "", WorkflowExecutionID: "exec-1"}, "0xabc", fee) + require.Empty(t, emitter.records) + require.Len(t, obs.FilterMessage("Gas usage meter record not emitted").All(), 1) + }) +} diff --git a/chain_capabilities/common/go.mod b/chain_capabilities/common/go.mod index 1481fe92d..f7fea57a2 100644 --- a/chain_capabilities/common/go.mod +++ b/chain_capabilities/common/go.mod @@ -5,10 +5,12 @@ go 1.27.1 require ( github.com/jpillora/backoff v1.0.0 github.com/smartcontractkit/capabilities/libs v0.0.0-20260714133332-db2a5f11cd64 - github.com/smartcontractkit/chainlink-common v0.11.2-0.20260917115705-1d3a14a9b049 + github.com/smartcontractkit/chainlink-common v0.11.2-0.20261006175240-29594528f464 + github.com/smartcontractkit/chainlink-protos/metering/go v0.0.0-20260710151514-27b5a126dabe github.com/smartcontractkit/libocr v0.0.0-20260810200708-618b5bf7f342 github.com/stretchr/testify v1.12.1 go.opentelemetry.io/otel v1.46.0 + go.uber.org/zap v1.28.0 google.golang.org/protobuf v1.36.12 ) @@ -32,6 +34,7 @@ require ( github.com/google/uuid v1.6.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0 // indirect github.com/invopop/jsonschema v0.13.0 // indirect + github.com/jonboulle/clockwork v0.5.0 // indirect github.com/json-iterator/go v1.1.12 // indirect github.com/leodido/go-urn v1.4.0 // indirect github.com/mailru/easyjson v0.9.0 // indirect @@ -49,7 +52,7 @@ require ( github.com/santhosh-tekuri/jsonschema/v5 v5.3.1 // indirect github.com/shopspring/decimal v1.4.0 // indirect github.com/smartcontractkit/chain-selectors v1.0.104 // indirect - github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92 // indirect + github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 // indirect github.com/smartcontractkit/chainlink-protos/cre/go v0.0.0-20260804191526-b7a850ae7648 // indirect github.com/smartcontractkit/chainlink-protos/linking-service/go v0.0.0-20251002192024-d2ad9222409b // indirect github.com/stretchr/objx v0.5.3 // indirect @@ -74,7 +77,6 @@ require ( go.opentelemetry.io/otel/trace v1.46.0 // indirect go.opentelemetry.io/proto/otlp v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect - go.uber.org/zap v1.28.0 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect go.yaml.in/yaml/v3 v3.0.5 // indirect golang.org/x/crypto v0.55.0 // indirect diff --git a/chain_capabilities/common/go.sum b/chain_capabilities/common/go.sum index a2c943415..2d6a64241 100644 --- a/chain_capabilities/common/go.sum +++ b/chain_capabilities/common/go.sum @@ -46,6 +46,8 @@ github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0 h1:/Tnpcb2E0Pz/tN9s3bfEY2Q8ePC github.com/grpc-ecosystem/grpc-gateway/v2 v2.30.0/go.mod h1:zOBXOsUaBSjKgmH4OGzV1esUpR3oUSCPYVd2cUBjKYY= github.com/invopop/jsonschema v0.13.0 h1:KvpoAJWEjR3uD9Kbm2HWJmqsEaHt8lBUpd0qHcIi21E= github.com/invopop/jsonschema v0.13.0/go.mod h1:ffZ5Km5SWWRAIN6wbDXItl95euhFz2uON45H2qjYt+0= +github.com/jonboulle/clockwork v0.5.0 h1:Hyh9A8u51kptdkR+cqRpT1EebBwTn1oK9YfGYbdFz6I= +github.com/jonboulle/clockwork v0.5.0/go.mod h1:3mZlmanh0g2NDKO5TWZVJAfofYk64M7XN3SzBPjZF60= github.com/jpillora/backoff v1.0.0 h1:uvFg412JmmHBHw7iwprIxkPMI+sGQ4kzOWsMeHnm2EA= github.com/jpillora/backoff v1.0.0/go.mod h1:J/6gKK9jxlEcS3zixgDgUAsiuZ7yrSoa/FX5e0EB2j4= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= @@ -92,14 +94,16 @@ github.com/smartcontractkit/capabilities/libs v0.0.0-20260714133332-db2a5f11cd64 github.com/smartcontractkit/capabilities/libs v0.0.0-20260714133332-db2a5f11cd64/go.mod h1:DZkOtlIRV6ecpCskdXMF+k2es3uJQqTT1m5TRB9etjQ= github.com/smartcontractkit/chain-selectors v1.0.104 h1:/n9pPGM5W/+r1eHoWZv4VwX9LNS1af4+ICyhM8zKRNM= github.com/smartcontractkit/chain-selectors v1.0.104/go.mod h1:qy7whtgG5g+7z0jt0nRyii9bLND9m15NZTzuQPkMZ5w= -github.com/smartcontractkit/chainlink-common v0.11.2-0.20260917115705-1d3a14a9b049 h1:WebE9HTubitMC5xq+xCAArEI+oBLekVOoxYnCkB8whI= -github.com/smartcontractkit/chainlink-common v0.11.2-0.20260917115705-1d3a14a9b049/go.mod h1:3CaCxSyZmgWl1PhxneFSZqOfN7DF2VnHML0bTfiDNIo= -github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92 h1:Z1OIGukDhJD+Gknr5G07Vtml24E7BwKaIQoWzV/SdqA= -github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20260915184316-2730f1867c92/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= +github.com/smartcontractkit/chainlink-common v0.11.2-0.20261006175240-29594528f464 h1:DVXwyIGekGqLaFAcSWa3K3S209gAY0WDFwliYdLp1dk= +github.com/smartcontractkit/chainlink-common v0.11.2-0.20261006175240-29594528f464/go.mod h1:37Mu0pbAnGmHPgOZ8bQvgP4hzXw/Wu7smKX4cN46UOs= +github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549 h1:J4vYRJxkBO8Yf//17wScKjYWmS45VZaFWx0zt+4ETEU= +github.com/smartcontractkit/chainlink-common/pkg/chipingress v0.0.11-0.20261005132539-b7d63d5fd549/go.mod h1:Scyk7O80fde5rJF3Ty3fEE3xiS6j84eWbu8/2gl3pIY= github.com/smartcontractkit/chainlink-protos/cre/go v0.0.0-20260804191526-b7a850ae7648 h1:WEUMkKQPAgcNMRgES6CBWrRUiII+HKEWQjulKQBSuMA= github.com/smartcontractkit/chainlink-protos/cre/go v0.0.0-20260804191526-b7a850ae7648/go.mod h1:/i8hjTPFdVWHiY+QjeSiVS2Z3GB3WAZznGgXHstC02E= github.com/smartcontractkit/chainlink-protos/linking-service/go v0.0.0-20251002192024-d2ad9222409b h1:QuI6SmQFK/zyUlVWEf0GMkiUYBPY4lssn26nKSd/bOM= github.com/smartcontractkit/chainlink-protos/linking-service/go v0.0.0-20251002192024-d2ad9222409b/go.mod h1:qSTSwX3cBP3FKQwQacdjArqv0g6QnukjV4XuzO6UyoY= +github.com/smartcontractkit/chainlink-protos/metering/go v0.0.0-20260710151514-27b5a126dabe h1:MDnY5wQbWTpFdDnMRicEnoMfSP5nM/KncARr4skP1ug= +github.com/smartcontractkit/chainlink-protos/metering/go v0.0.0-20260710151514-27b5a126dabe/go.mod h1:z7lx7wI3XZ4u9kmUtAVdwn1BCC9T8aieWSDcuDgPTdQ= github.com/smartcontractkit/libocr v0.0.0-20260810200708-618b5bf7f342 h1:pEcgcjTGA83MzpqbTbyIg9AJrOs62s77SooDdJGIg9w= github.com/smartcontractkit/libocr v0.0.0-20260810200708-618b5bf7f342/go.mod h1:5JPtsRwjugpyfsdEALC4RopfvohqK/G+3DHaR8uv+Bc= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=