Skip to content

Commit 05fbf40

Browse files
committed
CHASM telemetry
1 parent 44a3f6b commit 05fbf40

14 files changed

Lines changed: 190 additions & 69 deletions

File tree

chasm/registrable_task.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,13 @@ import (
44
"context"
55
"fmt"
66
"reflect"
7+
"strings"
78
)
89

910
type (
1011
RegistrableTask struct {
1112
taskType string
13+
telemetryTypeName string
1214
goType reflect.Type
1315
componentGoType reflect.Type // It is not clear how this one is used.
1416
validateFn validateFn
@@ -155,6 +157,7 @@ func (rt *RegistrableTask) registerToLibrary(
155157

156158
fqn := rt.fqType()
157159
rt.taskTypeID = GenerateTypeID(fqn)
160+
rt.telemetryTypeName = deriveTelemetryTypeName(rt.goType, fqn)
158161
return fqn, rt.taskTypeID, nil
159162
}
160163

@@ -173,3 +176,23 @@ func (rt *RegistrableTask) fqType() string {
173176
}
174177
return FullyQualifiedName(rt.library.Name(), rt.taskType)
175178
}
179+
180+
func (rt *RegistrableTask) telemetryType() string {
181+
if rt.telemetryTypeName != "" {
182+
return rt.telemetryTypeName
183+
}
184+
return rt.fqType()
185+
}
186+
187+
func deriveTelemetryTypeName(taskType reflect.Type, fallback string) string {
188+
for taskType.Kind() == reflect.Pointer {
189+
taskType = taskType.Elem()
190+
}
191+
// Remove common prefix/suffix from type name
192+
if name := taskType.Name(); name != "" {
193+
name = strings.TrimPrefix(name, "Chasm")
194+
name = strings.TrimSuffix(name, "Data")
195+
return name
196+
}
197+
return fallback
198+
}

chasm/registry.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -168,6 +168,15 @@ func (r *Registry) TaskFqnByID(id uint32) (string, bool) {
168168
return rt.fqType(), true
169169
}
170170

171+
// TaskTelemetryNameByID returns the human-readable telemetry label for a task type ID.
172+
func (r *Registry) TaskTelemetryNameByID(id uint32) (string, bool) {
173+
rt, ok := r.rtByID[id]
174+
if !ok {
175+
return "", false
176+
}
177+
return rt.telemetryType(), true
178+
}
179+
171180
// TaskIDFor converts registered task instance to task type ID.
172181
// This method should only be used by CHASM framework internal code,
173182
// NOT CHASM library developers.

chasm/registry_test.go

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -147,6 +147,18 @@ func (s *RegistryTestSuite) TestRegistry_RegisterTasks_Success() {
147147
require.True(s.T(), ok)
148148
require.Equal(s.T(), "TestLibrary.Task2", rt2.FqType())
149149

150+
taskID1, ok := r.TaskIDFor(testTask1{})
151+
require.True(s.T(), ok)
152+
name1, ok := r.TaskTelemetryNameByID(taskID1)
153+
require.True(s.T(), ok)
154+
require.Equal(s.T(), "testTask1", name1)
155+
156+
taskID2, ok := r.TaskIDFor(testTask2{})
157+
require.True(s.T(), ok)
158+
name2, ok := r.TaskTelemetryNameByID(taskID2)
159+
require.True(s.T(), ok)
160+
require.Equal(s.T(), "testTask2", name2)
161+
150162
tInstance2 := "invalid task instance"
151163
rt3, ok := r.TaskFor(tInstance2)
152164
require.False(s.T(), ok)

chasm/statemachine.go

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,10 @@ import (
44
"fmt"
55
"slices"
66

7+
"go.opentelemetry.io/otel/attribute"
8+
"go.opentelemetry.io/otel/trace"
79
"go.temporal.io/api/serviceerror"
10+
"go.temporal.io/server/common/telemetry"
811
)
912

1013
// ErrInvalidTransition is returned from [Transition.Apply] on an invalid state transition.
@@ -44,8 +47,24 @@ func (t Transition[S, SM, E]) Possible(sm SM) bool {
4447

4548
// Apply applies a transition event to the given state machine changing the state machine's state to the transition's
4649
// Destination on success.
47-
func (t Transition[S, SM, E]) Apply(sm SM, ctx MutableContext, event E) error {
50+
func (t Transition[S, SM, E]) Apply(sm SM, ctx MutableContext, event E) (retErr error) {
4851
prevState := sm.StateMachineState()
52+
53+
// Defer to always emit the transition telemetry event.
54+
if telemetry.DebugMode() {
55+
defer func() {
56+
attrs := []attribute.KeyValue{
57+
attribute.String("chasm.transition.source", fmt.Sprintf("%v", prevState)),
58+
attribute.String("chasm.transition.destination", fmt.Sprintf("%v", t.Destination)),
59+
}
60+
if retErr != nil {
61+
attrs = append(attrs, attribute.String("chasm.transition.error", retErr.Error()))
62+
}
63+
span := trace.SpanFromContext(ctx.goContext())
64+
span.AddEvent("chasm.transition", trace.WithAttributes(attrs...))
65+
}()
66+
}
67+
4968
if !t.Possible(sm) {
5069
return fmt.Errorf("%w from %v", ErrInvalidTransition, prevState)
5170
}

cmd/tools/genrpcserverinterceptors/main.go

Lines changed: 19 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -137,13 +137,6 @@ func workflowTagGetters(messageType reflect.Type, depth int) messageData {
137137
case messageType.AssignableTo(workflowExecutionGetterT):
138138
pd.WorkflowIDGetter = "GetWorkflowExecution().GetWorkflowId()"
139139
pd.RunIDGetter = "GetWorkflowExecution().GetRunId()"
140-
case messageType.AssignableTo(taskTokenGetterT):
141-
for _, ert := range excludeTaskTokenTypes {
142-
if messageType.AssignableTo(ert) {
143-
return pd
144-
}
145-
}
146-
pd.TaskTokenGetter = "GetTaskToken()"
147140
default:
148141
// Might have any combination of these, or none.
149142
if messageType.AssignableTo(workflowIDGetterT) {
@@ -160,6 +153,25 @@ func workflowTagGetters(messageType reflect.Type, depth int) messageData {
160153
}
161154
}
162155

156+
if pd.ActivityIDGetter == "" && messageType.AssignableTo(activityIDGetterT) {
157+
pd.ActivityIDGetter = "GetActivityId()"
158+
}
159+
if pd.OperationIDGetter == "" && messageType.AssignableTo(operationIDGetterT) {
160+
pd.OperationIDGetter = "GetOperationId()"
161+
}
162+
if messageType.AssignableTo(taskTokenGetterT) {
163+
excluded := false
164+
for _, ert := range excludeTaskTokenTypes {
165+
if messageType.AssignableTo(ert) {
166+
excluded = true
167+
break
168+
}
169+
}
170+
if !excluded {
171+
pd.TaskTokenGetter = "GetTaskToken()"
172+
}
173+
}
174+
163175
// Iterates over fields in order they defined in proto file, not proto index.
164176
// Order is important because the first match wins.
165177
for fieldNum := 0; fieldNum < messageType.Elem().NumField(); fieldNum++ {

cmd/tools/genrpcserverinterceptors/server_interceptors.tmpl

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,7 @@ func (wt *WorkflowTags) extractFrom{{.Server}}Message(message any) []tag.Tag {
1515
{{- range .Messages}}
1616
case {{.Type}}:
1717
{{- if or .TaskTokenGetter .WorkflowIDGetter .RunIDGetter .ActivityIDGetter .OperationIDGetter .ChasmRunIDGetter}}
18-
{{- if .TaskTokenGetter}}
19-
return wt.fromTaskToken(r.{{ .TaskTokenGetter}})
20-
{{- else}}
18+
{{- if or .WorkflowIDGetter .RunIDGetter .ActivityIDGetter .OperationIDGetter .ChasmRunIDGetter}}
2119
return []tag.Tag{
2220
{{if .WorkflowIDGetter}} tag.WorkflowID(r.{{.WorkflowIDGetter}}),
2321
{{end -}}
@@ -30,6 +28,8 @@ func (wt *WorkflowTags) extractFrom{{.Server}}Message(message any) []tag.Tag {
3028
{{if .ChasmRunIDGetter}} tag.ChasmRunID(r.{{.ChasmRunIDGetter}}),
3129
{{end -}}
3230
}
31+
{{- else if .TaskTokenGetter}}
32+
return wt.fromTaskToken(r.{{ .TaskTokenGetter}})
3333
{{- end}}
3434
{{- else}}
3535
return nil

common/log/tag/tags.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ import (
2121

2222
// LoggingCallAtKey is reserved tag
2323
const (
24+
ActivityIDKey = "activity-id"
25+
ChasmRunIDKey = "run-id"
2426
LoggingCallAtKey = "logging-call-at"
2527
WorkflowIDKey = "wf-id"
2628
WorkflowRunIDKey = "wf-run-id"
@@ -914,7 +916,7 @@ func ActivityInfo(activityInfo any) ZapTag {
914916

915917
// ActivityID returns tag for an activity ID
916918
func ActivityID(id string) ZapTag {
917-
return NewStringTag("activity-id", id)
919+
return NewStringTag(ActivityIDKey, id)
918920
}
919921

920922
// OperationID returns tag for a nexus operation ID
@@ -924,7 +926,7 @@ func OperationID(id string) ZapTag {
924926

925927
// ChasmRunID returns tag for an entity run ID
926928
func ChasmRunID(id string) ZapTag {
927-
return NewStringTag("run-id", id)
929+
return NewStringTag(ChasmRunIDKey, id)
928930
}
929931

930932
// ActivitySize returns a tag for a standalone activity size

common/rpc/interceptor/logtags/matching_service_server_gen.go

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

common/rpc/interceptor/logtags/workflow_service_server_gen.go

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

common/rpc/interceptor/logtags/workflow_tags.go

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,10 +54,22 @@ func (wt *WorkflowTags) fromTaskToken(taskTokenBytes []byte) []tag.Tag {
5454
if len(taskTokenBytes) == 0 {
5555
return nil
5656
}
57+
5758
taskToken, err := wt.serializer.Deserialize(taskTokenBytes)
5859
if err != nil {
5960
wt.logger.Warn("unable to deserialize task token while getting workflow tags", tag.Error(err))
6061
return nil
6162
}
62-
return []tag.Tag{tag.WorkflowID(taskToken.WorkflowId), tag.WorkflowRunID(taskToken.RunId)}
63+
64+
var tags []tag.Tag
65+
if taskToken.WorkflowId != "" {
66+
tags = append(tags, tag.WorkflowID(taskToken.WorkflowId))
67+
}
68+
if taskToken.RunId != "" {
69+
tags = append(tags, tag.WorkflowRunID(taskToken.RunId))
70+
}
71+
if len(taskToken.ComponentRef) > 0 && taskToken.ActivityId != "" {
72+
tags = append(tags, tag.ActivityID(taskToken.ActivityId))
73+
}
74+
return tags
6375
}

0 commit comments

Comments
 (0)