diff --git a/pkg/resourcemanager/resourcemanager.go b/pkg/resourcemanager/resourcemanager.go index 6a2f80d706..752867ad0d 100644 --- a/pkg/resourcemanager/resourcemanager.go +++ b/pkg/resourcemanager/resourcemanager.go @@ -52,6 +52,7 @@ package resourcemanager import ( "context" + "math/big" "strconv" "sync" "time" @@ -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) { + 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 diff --git a/pkg/resourcemanager/resourcemanager_test.go b/pkg/resourcemanager/resourcemanager_test.go index 824e2b2aae..8359377237 100644 --- a/pkg/resourcemanager/resourcemanager_test.go +++ b/pkg/resourcemanager/resourcemanager_test.go @@ -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. diff --git a/pkg/resourcemanager/workflow_usage.go b/pkg/resourcemanager/workflow_usage.go new file mode 100644 index 0000000000..105c7b4d2d --- /dev/null +++ b/pkg/resourcemanager/workflow_usage.go @@ -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:::", +// built from Utilization.ResourceId (":") 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:). 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) +} + +// WorkflowUsageResourceID returns the Utilization.ResourceId for a workflow +// execution: ":". Neither component may be empty or +// contain ':', since the consumer splits on the first ':'. +func WorkflowUsageResourceID(workflowID, executionID string) (string, error) { + if workflowID == "" || executionID == "" || strings.Contains(workflowID, ":") || strings.Contains(executionID, ":") { + return "", errWorkflowUsageResourceID + } + return workflowID + ":" + executionID, nil +} diff --git a/pkg/resourcemanager/workflow_usage_test.go b/pkg/resourcemanager/workflow_usage_test.go new file mode 100644 index 0000000000..d3eca29907 --- /dev/null +++ b/pkg/resourcemanager/workflow_usage_test.go @@ -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) +}