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
100 changes: 100 additions & 0 deletions chain_capabilities/common/gasmeter/gasmeter.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
// Package gasmeter emits the cre:workflow:gas:<chain_selector> 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 {
Comment thread
DylanTinianov marked this conversation as resolved.
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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: here and stop, I would remove it. I'm not in the camp of all public methods need comments. No value added. I won't start a religious war though so non-blocking.

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,
)
}
118 changes: 118 additions & 0 deletions chain_capabilities/common/gasmeter/gasmeter_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
8 changes: 5 additions & 3 deletions chain_capabilities/common/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
)

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down
12 changes: 8 additions & 4 deletions chain_capabilities/common/go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading