Skip to content
Closed
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
4 changes: 2 additions & 2 deletions core/scripts/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ require (
github.com/shopspring/decimal v1.4.0
github.com/smartcontractkit/chain-selectors v1.0.112
github.com/smartcontractkit/chainlink-ccip/chains/evm v0.0.0-20260918135944-fa1a268dac47
github.com/smartcontractkit/chainlink-common v0.11.2-0.20260930142354-8fdc7816bef3
github.com/smartcontractkit/chainlink-common v0.11.2-0.20261006175512-16997122620e
github.com/smartcontractkit/chainlink-common/keystore v1.3.1-0.20260903141829-ef07b52a737d
github.com/smartcontractkit/chainlink-deployments-framework v0.123.3
github.com/smartcontractkit/chainlink-evm v0.3.4-0.20261005112317-b723176adfe8
Expand Down Expand Up @@ -506,7 +506,7 @@ require (
github.com/smartcontractkit/chainlink-ccip/chains/solana/gobindings v0.0.0-20260916222901-720a003dab50 // indirect
github.com/smartcontractkit/chainlink-ccip/deployment v0.0.0-20260918135944-fa1a268dac47 // indirect
github.com/smartcontractkit/chainlink-ccv v0.13.1-0.20260918171034-c93b0d2ef3c0 // 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-confidential-compute v1.3.0 // indirect
github.com/smartcontractkit/chainlink-data-streams v1.1.2-0.20261002081259-6c2163b21db9 // indirect
github.com/smartcontractkit/chainlink-evm/contracts/cre/gobindings v0.0.0-20260403151002-2c91155b5501 // indirect
Expand Down
7 changes: 7 additions & 0 deletions core/scripts/go.sum

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

14 changes: 14 additions & 0 deletions core/services/cre/cre.go
Original file line number Diff line number Diff line change
Expand Up @@ -1023,6 +1023,20 @@ func newWorkflowRegistrySyncerV2(
handlerOpts = append(handlerOpts, syncerV2.WithSpecMeter(specMeter))
}

// Workflow compute usage records (gas comes from the chain-write plugins)
// share the [Metering].MeterRecordsEnabled gate with durable resource
// metering; there is no separate flag. Nil meter means engines emit nothing.
if meterRecordsEnabled {
usageRM := resourcemanager.NewResourceManager(lggr, resourcemanager.ResourceManagerConfig{
MeterRecordsEnabled: true,
Emitter: beholder.GetEmitter(),
})
usageIdentity := meterIdentity
usageIdentity.Service = resourcemanager.EmittingServiceWorkflowEngine
usageIdentity = resourcemanager.WithWorkflowUsagePool(usageIdentity, resourcemanager.ResourceTypeWorkflowCompute)
handlerOpts = append(handlerOpts, syncerV2.WithUsageMeter(usageRM, usageIdentity))
}

mc := capCfg.WorkflowRegistry().ModuleCache()
cacheEnabled := mc.Enabled()
diskMonitorEnabled := mc.DiskMonitorEnabled() || cacheEnabled
Expand Down
22 changes: 22 additions & 0 deletions core/services/workflows/syncer/v2/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,12 @@ type eventHandler struct {
// disabled: its handler-facing methods are nil-receiver safe no-ops.
specMeter *SpecMeter

// usageMeter emits per-capability workflow usage records (compute) from the
// engines this handler creates. Nil when [Metering].CapabilityUsageEnabled
// is off; engines then emit nothing.
usageMeter *resourcemanager.ResourceManager
usageIdentity resourcemanager.ResourceIdentity

orgResolver orgresolver.OrgResolver
secretsFetcher v2.SecretsFetcher
// localSecretOverrides is keyed by owner address; values are secret id -> secret value
Expand Down Expand Up @@ -198,6 +204,16 @@ func WithBillingClient(client metering.BillingClient) func(*eventHandler) {
// sub-service and reports storage transitions through EmitSpecDelta; all other
// metering concerns (ResourceManager lifecycle, identity, snapshots) live on
// the meter. A nil meter (metering disabled) is a valid no-op.
// WithUsageMeter supplies the ResourceManager and base identity used by engines
// to emit cre:workflow:compute usage MeterRecords. The handler owns the
// ResourceManager lifecycle as a sub-service.
func WithUsageMeter(rm *resourcemanager.ResourceManager, identity resourcemanager.ResourceIdentity) func(*eventHandler) {
return func(e *eventHandler) {
e.usageMeter = rm
e.usageIdentity = identity
}
}

func WithSpecMeter(sm *SpecMeter) func(*eventHandler) {
return func(e *eventHandler) {
e.specMeter = sm
Expand Down Expand Up @@ -417,6 +433,9 @@ func NewEventHandler(
if eh.triggerCoordinator != nil {
subs = append(subs, eh.triggerCoordinator)
}
if eh.usageMeter != nil {
subs = append(subs, eh.usageMeter)
}
return subs
},
Start: eh.start,
Expand Down Expand Up @@ -966,6 +985,7 @@ func (h *eventHandler) buildEngineConfig(ctx context.Context, workflowID, owner
}},
)
cfg := h.newV2EngineConfig(ctx, selectingModule, workflowID, owner, tag, sdkName, name, config)
cfg.ConfidentialExecutions = confidential

cfg.CachedTriggerSubscriptions = h.parseCachedTriggerSubscriptions(workflowID, cachedTriggerSubs)
h.wireTriggerSubscriptionCacheHook(cfg, workflowID)
Expand Down Expand Up @@ -1505,6 +1525,8 @@ func (h *eventHandler) newV2EngineConfig(
WorkflowRegistryAddress: h.workflowRegistryAddress,
WorkflowRegistryChainSelector: h.workflowRegistryChainSelector,
OrgResolver: h.orgResolver,
UsageMeter: h.usageMeter,
UsageIdentity: h.usageIdentity,
SecretsFetcher: h.secretsFetcher,
OverrideFetcher: h.overrideFetcherForOwner(owner),
DebugMode: h.debugMode,
Expand Down
43 changes: 43 additions & 0 deletions core/services/workflows/v2/base_engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
"github.com/smartcontractkit/chainlink-common/pkg/custmsg"
"github.com/smartcontractkit/chainlink-common/pkg/logger"
"github.com/smartcontractkit/chainlink-common/pkg/metrics"
"github.com/smartcontractkit/chainlink-common/pkg/resourcemanager"
"github.com/smartcontractkit/chainlink-common/pkg/services"
"github.com/smartcontractkit/chainlink-common/pkg/services/orgresolver"
"github.com/smartcontractkit/chainlink-common/pkg/settings/cresettings"
Expand Down Expand Up @@ -737,6 +738,12 @@ func (e *baseEngine) startExecution(ctx context.Context, event triggers.Coordina
"computeMs", computeDuration.Milliseconds())
}

// Capability usage record for ordinary compute. Emitted for successful and
// failed executions alike (the compute happened), outside the legacy
// metering block so a legacy metering failure cannot suppress it, and
// skipped when the enclave ran the execution (metered on that path).
e.emitComputeUsage(ctx, executionLogger, executionID, computeDuration)

if isMetering {
computeUnit := billing.ResourceType_name[int32(billing.ResourceType_RESOURCE_TYPE_COMPUTE)]
mrErr := meteringReport.Settle(computeUnit,
Expand Down Expand Up @@ -1081,3 +1088,39 @@ func resolveOrgID(ctx context.Context, resolver orgresolver.OrgResolver, workflo
}
return resolvedOrg{ID: orgID}
}

// emitComputeUsage emits the cre:workflow:compute usage MeterRecord for one
// execution and logs the emission. The log line is a contract consumed by the
// billing reconciler (fields: executionID, eventID, resourceType, value, orgID)
// and must stay stable. Fail-open: never returns an error to the execution.
func (e *baseEngine) emitComputeUsage(ctx context.Context, lggr logger.Logger, executionID string, computeDuration time.Duration) {
if e.cfg.UsageMeter == nil {
return
}
if e.cfg.ConfidentialExecutions != nil && e.cfg.ConfidentialExecutions.TookExecution(executionID) {
return
}
resourceID, err := resourcemanager.WorkflowUsageResourceID(e.cfg.WorkflowID, executionID)
if err != nil {
lggr.Errorw("Compute usage meter record not emitted", "err", err)
return
}
identity := e.cfg.UsageIdentity
if node := e.localNode.Load(); node != nil {
identity = identity.WithDonID(strconv.FormatUint(uint64(node.WorkflowDON.ID), 10))
}
valueMs := computeDuration.Milliseconds()
// The capability event id for compute is the execution id itself: one
// compute record per execution, identical on every node of the DON.
e.cfg.UsageMeter.EmitUsage(ctx, identity, executionID, valueMs, resourcemanager.UtilizationFields{
ResourceType: resourcemanager.ResourceTypeWorkflowCompute,
ResourceID: resourceID,
OrgID: e.orgID,
})
lggr.Infow("Emitted capability usage meter record",
"eventID", executionID,
"resourceType", resourcemanager.ResourceTypeWorkflowCompute,
"value", valueMs,
"orgID", e.orgID,
)
}
177 changes: 177 additions & 0 deletions core/services/workflows/v2/compute_usage_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
package v2_test

import (
"context"
"errors"
"sync"
"testing"

"github.com/stretchr/testify/mock"
"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"

Check failure on line 16 in core/services/workflows/v2/compute_usage_test.go

View workflow job for this annotation

GitHub Actions / GolangCI Lint (.)

File is not properly formatted (gci)
modulemocks "github.com/smartcontractkit/chainlink-common/pkg/workflows/wasm/host/mocks"
meteringpb "github.com/smartcontractkit/chainlink-protos/metering/go"

regmocks "github.com/smartcontractkit/chainlink-common/pkg/types/core/mocks"
capmocks "github.com/smartcontractkit/chainlink/v2/core/capabilities/mocks"
v2 "github.com/smartcontractkit/chainlink/v2/core/services/workflows/v2"
"github.com/smartcontractkit/chainlink/v2/core/utils/matches"
)

// recordingEmitter captures MeterRecords emitted by a ResourceManager.
type recordingEmitter struct {
mu sync.Mutex
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.mu.Lock()
defer r.mu.Unlock()
r.records = append(r.records, &rec)
return nil
}

func (r *recordingEmitter) all() []*meteringpb.MeterRecord {
r.mu.Lock()
defer r.mu.Unlock()
return append([]*meteringpb.MeterRecord(nil), r.records...)
}

type staticTracker struct{ confidential bool }

func (s staticTracker) TookExecution(string) bool { return s.confidential }

// runOneExecution starts an engine with cfg, fires one trigger event whose
// module execution returns execErr, and returns the execution id.
func runOneExecution(t *testing.T, cfg *v2.EngineConfig, execErr error) string {
module := modulemocks.NewModuleV2(t)
module.EXPECT().Start()
module.EXPECT().Close()
capreg := regmocks.NewCapabilitiesRegistry(t)
capreg.EXPECT().LocalNode(matches.AnyContext).Return(newNode(t), nil)

initDoneCh := make(chan error)
subscribedCh := make(chan []string, 1)
finishedCh := make(chan string)
cfg.Module = module
cfg.CapRegistry = capreg
cfg.BillingClient = setupMockBillingClient(t)
cfg.OrgResolver = &mockOrgResolver{orgID: "org-123"}
cfg.Hooks = v2.LifecycleHooks{
OnInitialized: func(err error) { initDoneCh <- err },
OnSubscribedToTriggers: func(ids []string) { subscribedCh <- ids },
OnExecutionFinished: func(id string, _ string) { finishedCh <- id },
}

engine, err := v2.NewEngine(cfg)
require.NoError(t, err)

module.EXPECT().Execute(matches.AnyContext, mock.Anything, mock.Anything).Return(newTriggerSubs(1), nil).Once()
trigger := capmocks.NewTriggerCapability(t)
capreg.EXPECT().GetTrigger(matches.AnyContext, "id_0").Return(trigger, nil).Once()
eventCh := make(chan capabilities.TriggerResponse)
trigger.EXPECT().RegisterTrigger(matches.AnyContext, mock.Anything).Return(eventCh, nil).Once()
trigger.EXPECT().UnregisterTrigger(matches.AnyContext, mock.Anything).Return(nil).Once()
trigger.EXPECT().AckEvent(matches.AnyContext, mock.Anything, mock.Anything, mock.Anything).Return(nil)

require.NoError(t, engine.Start(t.Context()))
require.NoError(t, <-initDoneCh)
require.Equal(t, []string{"id_0"}, <-subscribedCh)

module.EXPECT().Execute(matches.AnyContext, mock.Anything, mock.Anything).Return(nil, execErr).Once()
eventCh <- capabilities.TriggerResponse{Event: capabilities.TriggerEvent{TriggerType: "basic-trigger@1.0.0", ID: "usage_event"}}
executionID := <-finishedCh
require.NoError(t, engine.Close())
return executionID
}

func TestEngine_ComputeUsageMeterRecord(t *testing.T) {
t.Parallel()

newUsageMeter := func(t *testing.T) (*resourcemanager.ResourceManager, *recordingEmitter) {
emitter := &recordingEmitter{}
rm := resourcemanager.NewResourceManager(logger.Test(t), resourcemanager.ResourceManagerConfig{
MeterRecordsEnabled: true,
Emitter: emitter,
})
return rm, emitter
}
identity := resourcemanager.WithWorkflowUsagePool(resourcemanager.ResourceIdentity{Product: "cre", Service: resourcemanager.EmittingServiceWorkflowEngine}, resourcemanager.ResourceTypeWorkflowCompute)

t.Run("emits one compute record per execution and logs the contract line", func(t *testing.T) {

Check failure on line 110 in core/services/workflows/v2/compute_usage_test.go

View workflow job for this annotation

GitHub Actions / GolangCI Lint (.)

Function TestEngine_ComputeUsageMeterRecord missing the call to method parallel in the test run (paralleltest)
lggr, obs := logger.TestObserved(t, zapcore.InfoLevel)
cfg := defaultTestConfig(t, nil)
cfg.Lggr = lggr
rm, emitter := newUsageMeter(t)
cfg.UsageMeter = rm
cfg.UsageIdentity = identity

executionID := runOneExecution(t, cfg, nil)

records := emitter.all()
require.Len(t, records, 1)
rec := records[0]
require.Equal(t, meteringpb.MeterAction_METER_ACTION_USAGE, rec.GetAction())
require.Equal(t, resourcemanager.EmittingServiceWorkflowEngine, rec.GetIdentity().GetService())
require.Equal(t, "cre:workflow:compute", rec.GetIdentity().GetResourcePool())
require.Equal(t, "cre:workflow:compute", rec.GetIdentity().GetResourcePoolId())
require.NotEmpty(t, rec.GetIdentity().GetDon().GetDonId(), "don id must be stamped from the local node")
require.Len(t, rec.GetUtilizations(), 1)
u := rec.GetUtilizations()[0]
require.Equal(t, resourcemanager.ResourceTypeWorkflowCompute, u.GetResourceType())
require.Equal(t, cfg.WorkflowID+":"+executionID, u.GetResourceId())
require.Equal(t, executionID, u.GetEventId(), "compute capability event id is the execution id")
require.Equal(t, "org-123", u.GetOrgId())
require.NotEmpty(t, u.GetValue())

logs := obs.FilterMessage("Emitted capability usage meter record").All()
require.Len(t, logs, 1)
fields := logs[0].ContextMap()
require.Equal(t, executionID, fields["eventID"])
require.Equal(t, resourcemanager.ResourceTypeWorkflowCompute, fields["resourceType"])
require.Equal(t, "org-123", fields["orgID"])
require.Contains(t, fields, "value")
})

t.Run("emits for failed executions too", func(t *testing.T) {

Check failure on line 145 in core/services/workflows/v2/compute_usage_test.go

View workflow job for this annotation

GitHub Actions / GolangCI Lint (.)

Function TestEngine_ComputeUsageMeterRecord missing the call to method parallel in the test run (paralleltest)
cfg := defaultTestConfig(t, nil)
rm, emitter := newUsageMeter(t)
cfg.UsageMeter = rm
cfg.UsageIdentity = identity

runOneExecution(t, cfg, errors.New("module failed"))

require.Len(t, emitter.all(), 1)
})

t.Run("skips executions delegated to the confidential module", func(t *testing.T) {

Check failure on line 156 in core/services/workflows/v2/compute_usage_test.go

View workflow job for this annotation

GitHub Actions / GolangCI Lint (.)

Function TestEngine_ComputeUsageMeterRecord missing the call to method parallel in the test run (paralleltest)
cfg := defaultTestConfig(t, nil)
rm, emitter := newUsageMeter(t)
cfg.UsageMeter = rm
cfg.UsageIdentity = identity
cfg.ConfidentialExecutions = staticTracker{confidential: true}

runOneExecution(t, cfg, nil)

require.Empty(t, emitter.all())
})

t.Run("emits nothing when no usage meter is configured", func(t *testing.T) {

Check failure on line 168 in core/services/workflows/v2/compute_usage_test.go

View workflow job for this annotation

GitHub Actions / GolangCI Lint (.)

Function TestEngine_ComputeUsageMeterRecord missing the call to method parallel in the test run (paralleltest)
lggr, obs := logger.TestObserved(t, zapcore.InfoLevel)
cfg := defaultTestConfig(t, nil)
cfg.Lggr = lggr

runOneExecution(t, cfg, nil)

require.Empty(t, obs.FilterMessage("Emitted capability usage meter record").All())
})
}
16 changes: 16 additions & 0 deletions core/services/workflows/v2/confidential_module.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,8 @@ func ParseWorkflowAttributes(data []byte) (WorkflowAttributes, error) {
// Instead of running WASM locally, it delegates execution to the
// confidential-workflows capability via the CapabilitiesRegistry.
type ConfidentialModule struct {
// handled marks executions delegated to the enclave; see TookExecution.
handled sync.Map
capRegistry registry.CapabilitiesRegistry
binaryURL string
binaryHash []byte
Expand Down Expand Up @@ -169,6 +171,7 @@ func (m *ConfidentialModule) Execute(
}
m.executionHandlers.AddExecution(m.workflowID, workflowExecutionID, rawSecretsHelper)
defer m.executionHandlers.RemoveExecution(m.workflowID, workflowExecutionID)
m.handled.Store(workflowExecutionID, struct{}{})

requirements := loadAndDelete[*sdkpb.Requirements](&m.requirements, workflowExecutionID)
restrictions := loadAndDelete[*sdkpb.Restrictions](&m.restritions, workflowExecutionID)
Expand Down Expand Up @@ -211,6 +214,19 @@ func (m *ConfidentialModule) Execute(
return capOutput.SdkExecutionResult, nil
}

// ConfidentialExecutionTracker answers whether an execution ran in the enclave.
// TookExecution consumes the mark, so it must be called at most once per
// execution, after the module's Execute has returned.
type ConfidentialExecutionTracker interface {
TookExecution(executionID string) bool
}

// TookExecution implements ConfidentialExecutionTracker.
func (m *ConfidentialModule) TookExecution(executionID string) bool {
_, ok := m.handled.LoadAndDelete(executionID)
return ok
}

func (m *ConfidentialModule) SetRequirements(executionID string, requirements *sdkpb.Requirements) {
m.requirements.Store(executionID, requirements)
}
Expand Down
Loading
Loading