Repository navigation
Billing: Shared gasmeter for cre:workflow:gas usage records #813
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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 { | ||
| 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. | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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, | ||
| ) | ||
| } | ||
| 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) | ||
| }) | ||
| } |
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Uh oh!
There was an error while loading. Please reload this page.