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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ to docs, or any other relevant information.
workflow flow types to activities, adds flow/step/subflow/RPC metric tags, and rewrites Temporal
SDK metric names into the `dex_*` namespace.

### Changed

- Activities with a Dex `FlowTypeProvider` no longer store the parent workflow flow type in their
headers. The provider is authoritative, and an empty result produces the `none` metric label.

### Breaking Changes

- SDK metrics beginning with `temporal_` are now emitted under their corresponding `dex_*` names.
Expand Down
5 changes: 2 additions & 3 deletions internal/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,9 +103,8 @@ type (
// invoking the activity function. Empty values become "none". Provider
// results are metric labels and should have bounded cardinality.
//
// FlowTypeProvider extracts the flow type. A non-empty value overrides the
// flow type inherited from the parent workflow; an empty value falls back to
// the inherited value, or "none" when no inherited value exists.
// FlowTypeProvider extracts the flow type instead of inheriting it from the
// parent workflow. An empty value becomes "none".
FlowTypeProvider func(input any) string
// StepTypeProvider identifies the activity as a step and extracts its step
// type. At most one of StepTypeProvider, SubFlowTypeProvider, and
Expand Down
17 changes: 11 additions & 6 deletions internal/dex_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,18 +117,16 @@ func (p dexActivityMetricProviders) values(typeName string, input any, inherited
rpcName: metrics.NoneTagValue,
kind: p.metricKind(),
}
if inheritedFlowType != "" {
if p.flowTypeProvider == nil && inheritedFlowType != "" {
values.flowType = inheritedFlowType
}
var err error
if p.flowTypeProvider != nil {
flowType, providerErr := invokeDexMetricsProviderRaw("activity", typeName, "FlowTypeProvider", p.flowTypeProvider, input)
flowType, providerErr := invokeDexMetricsProvider("activity", typeName, "FlowTypeProvider", p.flowTypeProvider, input)
if providerErr != nil {
return values, providerErr
}
if flowType != "" {
values.flowType = flowType
}
values.flowType = flowType
}
switch p.kind {
case dexActivityMetricKindStep:
Expand All @@ -155,7 +153,14 @@ type dexWorkflowMetricsEnvironment interface {
dexWorkflowFlowType() (string, bool)
}

func addDexFlowTypeHeader(header *commonpb.Header, env WorkflowEnvironment) {
func addDexFlowTypeHeader(
header *commonpb.Header,
env WorkflowEnvironment,
activityProviders dexActivityMetricProviders,
) {
if activityProviders.flowTypeProvider != nil {
return
}
dexEnv, ok := env.(dexWorkflowMetricsEnvironment)
if !ok {
return
Expand Down
35 changes: 28 additions & 7 deletions internal/dex_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,13 +46,16 @@ func TestDexActivityMetricsProviders(t *testing.T) {

values, err = providers.values("activity", &dexMetricsTestInput{}, "parent")
require.NoError(t, err)
require.Equal(t, "parent", values.flowType)
require.Equal(t, "none", values.flowType)
require.Equal(t, "none", values.stepType)

providers.flowTypeProvider = func(any) string { return "none" }
values, err = providers.values("activity", &dexMetricsTestInput{}, "parent")
inheritedProviders := dexActivityProviders(RegisterActivityOptions{
StepTypeProvider: func(input any) string { return input.(*dexMetricsTestInput).StepType },
})
values, err = inheritedProviders.values("activity", &dexMetricsTestInput{StepType: "step"}, "parent")
require.NoError(t, err)
require.Equal(t, "none", values.flowType)
require.Equal(t, "parent", values.flowType)
require.Equal(t, "step", values.stepType)
}

func TestDexMetricsProviderPanicReturnsError(t *testing.T) {
Expand Down Expand Up @@ -86,11 +89,25 @@ func TestDexFlowTypeHeader(t *testing.T) {
require.Empty(t, popDexFlowTypeHeader(nil))
require.Empty(t, popDexFlowTypeHeader(&commonpb.Header{}))

header := &commonpb.Header{Fields: map[string]*commonpb.Payload{}}
env := &workflowEnvironmentImpl{dexFlowType: "flow", dexFlowTypeConfigured: true}
addDexFlowTypeHeader(header, env)
header := &commonpb.Header{Fields: map[string]*commonpb.Payload{}}
addDexFlowTypeHeader(header, env, dexActivityMetricProviders{})
require.Equal(t, "flow", popDexFlowTypeHeader(header))
require.NotContains(t, header.Fields, dexFlowTypeHeaderName)

header = &commonpb.Header{Fields: map[string]*commonpb.Payload{}}
providers := dexActivityProviders(RegisterActivityOptions{
FlowTypeProvider: func(input any) string { return input.(*dexMetricsTestInput).FlowType },
})
addDexFlowTypeHeader(header, env, providers)
require.NotContains(t, header.Fields, dexFlowTypeHeaderName)

header = &commonpb.Header{Fields: map[string]*commonpb.Payload{}}
providers = dexActivityProviders(RegisterActivityOptions{
StepTypeProvider: func(input any) string { return input.(*dexMetricsTestInput).StepType },
})
addDexFlowTypeHeader(header, env, providers)
require.Equal(t, "flow", popDexFlowTypeHeader(header))
}

func TestDexActivityProviderUsesDecodedFirstArgumentOnce(t *testing.T) {
Expand Down Expand Up @@ -177,7 +194,11 @@ func TestDexWorkflowAndActivityMetricsPropagation(t *testing.T) {
},
})

input := &dexMetricsTestInput{FlowType: "OrderFlow", StepType: "ChargeCard"}
input := &dexMetricsTestInput{
FlowType: "OrderFlow",
ActivityFlowType: "OrderFlow",
StepType: "ChargeCard",
}
env.ExecuteWorkflow("dexMetricsWorkflow", input)
require.True(t, env.IsWorkflowCompleted())
require.NoError(t, env.GetWorkflowError())
Expand Down
8 changes: 5 additions & 3 deletions internal/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -1083,7 +1083,8 @@ func (wc *workflowEnvironmentInterceptor) ExecuteActivity(ctx Context, typeName
settable.Set(nil, err)
return future
}
addDexFlowTypeHeader(header, getWorkflowEnvironment(ctx))
activityProviders := registry.getDexActivityMetricsProviders(typeName)
addDexFlowTypeHeader(header, getWorkflowEnvironment(ctx), activityProviders)

env := getWorkflowEnvironment(ctx)
// Generate activity ID before serialization so it's available to context-aware data converters
Expand Down Expand Up @@ -1214,14 +1215,15 @@ func ExecuteLocalActivity(ctx Context, activity interface{}, args ...interface{}

func (wc *workflowEnvironmentInterceptor) ExecuteLocalActivity(ctx Context, typeName string, args ...interface{}) Future {
future, settable := newDecodeFuture(ctx, typeName)
activityProviders := getRegistryFromWorkflowContext(ctx).getDexActivityMetricsProviders(typeName)

envOptions := getWorkflowEnvOptions(ctx)
header, err := workflowHeaderPropagated(ctx, envOptions.ContextPropagators)
if err != nil {
settable.Set(nil, err)
return future
}
addDexFlowTypeHeader(header, getWorkflowEnvironment(ctx))
addDexFlowTypeHeader(header, getWorkflowEnvironment(ctx), activityProviders)

var activityFn interface{}
localCtx := ctx.Value(localActivityFnContextKey).(*localActivityContext)
Expand Down Expand Up @@ -1299,7 +1301,7 @@ func (wc *workflowEnvironmentInterceptor) ExecuteLocalActivity(ctx Context, type
ScheduledTime: Now(ctx), // initial scheduled time
Header: header,
Attempt: 1, // Attempts always start at one
DexMetricsProviders: getRegistryFromWorkflowContext(ctx).getDexActivityMetricsProviders(typeName),
DexMetricsProviders: activityProviders,
}

Go(ctx, func(ctx Context) {
Expand Down
49 changes: 49 additions & 0 deletions test/dex_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (

"github.com/google/uuid"
"github.com/stretchr/testify/require"
enumspb "go.temporal.io/api/enums/v1"
"go.temporal.io/sdk/activity"
"go.temporal.io/sdk/client"
"go.temporal.io/sdk/temporal"
Expand All @@ -16,6 +17,7 @@ import (
)

const (
dexFlowTypeHeaderName = "__temporal_sdk_dex_flow_type"
dexMetricsWorkflowName = "dex-metrics-workflow"
dexSystemWorkflowName = "dex-system-workflow"
dexSyncStepActivityName = "dex-sync-step-activity"
Expand Down Expand Up @@ -148,6 +150,53 @@ func (ts *IntegrationTestSuite) TestDexMetricsProvidersAndNames() {
}
}

func (ts *IntegrationTestSuite) TestDexFlowTypeHeaderInheritance() {
input := &dexMetricsIntegrationInput{
FlowType: "OrderFlow",
StepType: "ChargeCard",
SubFlowType: "Fulfillment",
RPCName: "ReserveInventory",
}
run, err := ts.client.ExecuteWorkflow(context.Background(), client.StartWorkflowOptions{
ID: "dex-flow-type-header-" + uuid.NewString(),
TaskQueue: ts.taskQueueName,
}, dexMetricsWorkflowName, input)
ts.NoError(err)
ts.NoError(run.Get(context.Background(), nil))

expectedHeaderByActivityType := map[string]bool{
dexSyncStepActivityName: false,
dexSubFlowActivityName: true,
dexRPCActivityName: false,
dexSystemActivityName: true,
}
observedHeaderByActivityType := make(map[string]bool, len(expectedHeaderByActivityType))
history := ts.client.GetWorkflowHistory(
context.Background(),
run.GetID(),
run.GetRunID(),
false,
enumspb.HISTORY_EVENT_FILTER_TYPE_ALL_EVENT,
)
for history.HasNext() {
event, historyErr := history.Next()
ts.NoError(historyErr)
attributes := event.GetActivityTaskScheduledEventAttributes()
if attributes == nil {
continue
}
activityType := attributes.GetActivityType().GetName()
expectedHeader, tracked := expectedHeaderByActivityType[activityType]
if !tracked {
continue
}
_, hasHeader := attributes.GetHeader().GetFields()[dexFlowTypeHeaderName]
ts.Equal(expectedHeader, hasHeader, activityType)
observedHeaderByActivityType[activityType] = hasHeader
}
ts.Equal(expectedHeaderByActivityType, observedHeaderByActivityType)
}

func (ts *IntegrationTestSuite) TestDexMetricsProviderPanicsSkipBusinessFunctions() {
run, err := ts.client.ExecuteWorkflow(context.Background(), client.StartWorkflowOptions{
ID: "dex-workflow-provider-panic-" + uuid.NewString(),
Expand Down
Loading