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
8 changes: 8 additions & 0 deletions pkg/resourcemanager/resourcemanager.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ package resourcemanager

import (
"context"
"math/big"
"strconv"
"sync"
"time"
Expand Down Expand Up @@ -444,6 +445,13 @@ func (rm *ResourceManager) EmitUsage(ctx context.Context, identity ResourceIdent
rm.emitRecord(ctx, identity, eventID, meteringpb.MeterAction_METER_ACTION_USAGE, NewUtilizationInt(quantity, fields))
}

// EmitUsageValue is EmitUsage for quantities that do not fit an int64, such as a
// transaction fee in wei. A nil value emits "0". Same event_id contract and
// fail-open semantics as EmitUsage.
func (rm *ResourceManager) EmitUsageValue(ctx context.Context, identity ResourceIdentity, eventID string, value *big.Int, fields UtilizationFields) {
Comment thread
DylanTinianov marked this conversation as resolved.

@patrickhuie19 patrickhuie19 Oct 6, 2026 •

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.

one note is that I'm not sure how the *big.Int will go across the wire in the message payload. it will need to be serialized to string, we should double check that the emitRecord call will handle this properly. we use a cloudEvent packaging, which I believe wraps OTEL

rm.emitRecord(ctx, identity, eventID, meteringpb.MeterAction_METER_ACTION_USAGE, NewUtilizationBig(value, fields))
}

// emitRecord builds and emits a single-utilization MeterRecord for action.
// eventID is supplied by the caller and stamped on the utilization; it is never
// generated here. An empty eventID means the producer failed to derive a
Expand Down
29 changes: 29 additions & 0 deletions pkg/resourcemanager/resourcemanager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,35 @@ func TestEmitUsage_Action(t *testing.T) {
assert.Equal(t, "usage:req-1", record.GetUtilizations()[0].GetEventId())
}

func TestEmitUsageValue_Action(t *testing.T) {
emitter := &fakeEmitter{}
rm := NewResourceManager(logger.Test(t), ResourceManagerConfig{
MeterRecordsEnabled: true,
Emitter: emitter,
})

fee, ok := new(big.Int).SetString("123456789012345678901234567890", 10) // > MaxInt64
require.True(t, ok)
rm.EmitUsageValue(t.Context(), testIdentity, "0xabc", fee, UtilizationFields{
ResourceType: WorkflowGasResourceType(1),
ResourceID: "wf-1:exec-1",
OrgID: "org-1",
})
rm.EmitUsageValue(t.Context(), testIdentity, "0xdef", nil, UtilizationFields{ResourceType: WorkflowGasResourceType(1), ResourceID: "wf-1:exec-2"})

require.Len(t, emitter.calls, 2)
var record meteringpb.MeterRecord
require.NoError(t, proto.Unmarshal(emitter.calls[0].body, &record))
assert.Equal(t, meteringpb.MeterAction_METER_ACTION_USAGE, record.GetAction())
require.Len(t, record.GetUtilizations(), 1)
assert.Equal(t, "123456789012345678901234567890", record.GetUtilizations()[0].GetValue())
assert.Equal(t, "cre:workflow:gas:1", record.GetUtilizations()[0].GetResourceType())
assert.Equal(t, "0xabc", record.GetUtilizations()[0].GetEventId())

require.NoError(t, proto.Unmarshal(emitter.calls[1].body, &record))
assert.Equal(t, "0", record.GetUtilizations()[0].GetValue())
}

// TestEmitDelta_EventIDIdenticalAcrossNodes proves the core cross-node contract:
// two nodes fielding the SAME logical delta (identical eventID) emit the
// identical event_id, while a distinct request (distinct eventID) is distinct.
Expand Down
77 changes: 77 additions & 0 deletions pkg/resourcemanager/workflow_usage.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package resourcemanager

import (
"errors"
"strconv"
"strings"
)

// Workflow capability usage records: one METER_ACTION_USAGE MeterRecord per
// billable capability event of a workflow execution (compute duration, gas per
// chain write). The billing service turns each into a billing_records row with
// event id "cre:workflow:<workflow_id>:<execution_id>:<capability_event_id>",
// built from Utilization.ResourceId ("<workflow_id>:<execution_id>") and
// Utilization.EventId (the capability event id). Both must be identical on every
// node of the DON: they are the quorum and dedup key.
const (
// WorkflowRecordType is the billing record type shared by all workflow
// capability usage resource types.
WorkflowRecordType = "cre:workflow"

// ResourceTypeWorkflowCompute meters ordinary workflow compute in
// milliseconds. The capability event id is the workflow execution id.
ResourceTypeWorkflowCompute = WorkflowRecordType + ":compute"

// ResourceTypeWorkflowGasPrefix is the prefix of the per-chain gas resource
// type; see WorkflowGasResourceType. Values are in the chain's native
// smallest unit (wei, lamports). The capability event id is the tx hash.
ResourceTypeWorkflowGasPrefix = WorkflowRecordType + ":gas:"

// EmittingServiceChainWrite is the Identity.Service used by chain-write
// capabilities for gas usage records.
EmittingServiceChainWrite = "chain-write"
// EmittingServiceWorkflowEngine is the Identity.Service used by the
// workflow engine for compute usage records.
EmittingServiceWorkflowEngine = "workflow-engine"

// WorkflowGasResourcePool is the Identity.ResourcePool for gas usage
// records; ResourcePoolID is the fully qualified resource type
// (cre:workflow:gas:<chain_selector>). See WithWorkflowUsagePool.
WorkflowGasResourcePool = WorkflowRecordType + ":gas"
// WorkflowComputeResourcePool is the Identity.ResourcePool for compute
// usage records; ResourcePoolID is the same value.
WorkflowComputeResourcePool = ResourceTypeWorkflowCompute
)

var errWorkflowUsageResourceID = errors.New("workflow usage resource id: workflow id and execution id must be non-empty and contain no ':'")

// WithWorkflowUsagePool returns id with ResourcePool and ResourcePoolID set
// for a workflow usage resource type, per the billing payload contract:
// pool is the resource type without its chain selector
// ("cre:workflow:gas", "cre:workflow:compute") and pool id is the fully
// qualified resource type.
func WithWorkflowUsagePool(id ResourceIdentity, resourceType string) ResourceIdentity {
if strings.HasPrefix(resourceType, ResourceTypeWorkflowGasPrefix) {
id.ResourcePool = WorkflowGasResourcePool
} else {
id.ResourcePool = resourceType
}
id.ResourcePoolID = resourceType
return id
}

// WorkflowGasResourceType returns the gas resource type for a chain selector,
// e.g. "cre:workflow:gas:421614".
func WorkflowGasResourceType(chainSelector uint64) string {
return ResourceTypeWorkflowGasPrefix + strconv.FormatUint(chainSelector, 10)

@patrickhuie19 patrickhuie19 Oct 6, 2026 •

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Matches that section: ResourcePool is cre:workflow:gas / cre:workflow:compute, ResourcePoolId is the fully qualified resource type, ResourceId is <wf>:<exec>.

}

// WorkflowUsageResourceID returns the Utilization.ResourceId for a workflow
// execution: "<workflow_id>:<execution_id>". Neither component may be empty or
// contain ':', since the consumer splits on the first ':'.
func WorkflowUsageResourceID(workflowID, executionID string) (string, error) {

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: is this shared with any other components? maybe deserves its own file or even a "utils" package?

if workflowID == "" || executionID == "" || strings.Contains(workflowID, ":") || strings.Contains(executionID, ":") {
return "", errWorkflowUsageResourceID
}
return workflowID + ":" + executionID, nil
}
36 changes: 36 additions & 0 deletions pkg/resourcemanager/workflow_usage_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
package resourcemanager

import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestWorkflowGasResourceType(t *testing.T) {
assert.Equal(t, "cre:workflow:gas:421614", WorkflowGasResourceType(421614))
assert.Equal(t, "cre:workflow:compute", ResourceTypeWorkflowCompute)
}

func TestWorkflowUsageResourceID(t *testing.T) {
id, err := WorkflowUsageResourceID("wf-1", "exec-1")
require.NoError(t, err)
assert.Equal(t, "wf-1:exec-1", id)

for _, tc := range [][2]string{{"", "exec-1"}, {"wf-1", ""}, {"wf:1", "exec-1"}, {"wf-1", "exec:1"}} {
_, err := WorkflowUsageResourceID(tc[0], tc[1])
assert.ErrorIs(t, err, errWorkflowUsageResourceID, "%q %q", tc[0], tc[1])
}
}

func TestWithWorkflowUsagePool(t *testing.T) {
base := ResourceIdentity{Product: "cre", Service: EmittingServiceChainWrite}
gas := WithWorkflowUsagePool(base, WorkflowGasResourceType(421614))
assert.Equal(t, "cre:workflow:gas", gas.ResourcePool)
assert.Equal(t, "cre:workflow:gas:421614", gas.ResourcePoolID)
assert.Equal(t, EmittingServiceChainWrite, gas.Service, "other fields untouched")

compute := WithWorkflowUsagePool(base, ResourceTypeWorkflowCompute)
assert.Equal(t, "cre:workflow:compute", compute.ResourcePool)
assert.Equal(t, "cre:workflow:compute", compute.ResourcePoolID)
}
Loading