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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Changed

- Activities, sub-orchestrations, and ContinueAsNew do not add owning-orchestration identity tags, aligning automatic tags with .NET. Activity work items supply the owning instance ID and the activity's own identity; orchestrator contexts use native history fields. User-tag inheritance, immutable context fields, namespace markers, and distributed tracing are preserved.
- DTS token audiences now default per options/client/worker instance to `https://durabletask.azure.us` when `REGION_NAME` starts with `usgov` or `usdod` (case-insensitively), otherwise `https://durabletask.io`. Explicit `Options.ResourceID` or connection-string `ResourceId` overrides the default without changing the endpoint or credential authority. Resource audiences normalize surrounding whitespace, trailing slashes, and one existing `/.default` suffix before token requests; values that normalize to empty are rejected. Audiences remain stable across token refreshes and reconnects. Government-region applications requiring the old audience must explicitly set `https://durabletask.io`.
- Restore DTS-provided activity trace parents on `ActivityContext.Context()` without emitting duplicate SDK durable spans, preserving trace continuity for application instrumentation.
- Missing-instance orchestration waits now return `api.ErrInstanceNotFound` immediately instead of retrying `NotFound` until the caller deadline.
Expand Down
13 changes: 13 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -517,6 +517,19 @@ The SDK does not save the identity of the converter. A new converter must contin

Use `api.WithTags`, `task.WithActivityTags`, and `task.WithSubOrchestrationTags` to attach user tags. An activity and a sub-orchestration inherit the tags of the parent orchestration. A tag on the action has priority over an inherited tag. The completion actions carry the current tags, so ContinueAsNew keeps them.

Activities, sub-orchestrations, and ContinueAsNew do not add owning-orchestration
identity tags. Without caller tags or context fields, their tag maps
are absent on the wire, matching .NET's default automatic-tag behavior. Go keeps
its user-tag inheritance and immutable context-field propagation; when either is
present, the encoding marker separates user tags from context fields.
Distributed trace context remains a separate protocol field.

Activity contexts expose the activity's own name, version, and task ID in
`api.ActivityContextInfo`, and the owning instance ID in
`api.OrchestrationContextInfo`. Orchestrator contexts obtain their full identity
from native history fields, including after sub-orchestration starts and
ContinueAsNew version changes.

The client sends the sampled caller trace context when it schedules an orchestration or signals an entity. The worker adds separate action trace contexts for the service-owned activity and sub-orchestration spans. The worker does not emit duplicate local Durable Task spans. V2 entity requests do not carry per-operation trace context, so entity-emitted actions cannot inherit it.

Use `task.OrchestrationOptions.MaxEventsPerTurn` to limit the new events in one turn. If the worker uses only part of a batch, it sets `numEventsProcessed`. DTS then keeps the remaining events for the next replay. This count obeys the DTS work-item rules. The orchestration control markers do not count against the limit.
Expand Down
4 changes: 3 additions & 1 deletion api/context.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,15 @@ import (
"maps"
)

// ReservedContextFieldPrefix is reserved for Durable Task runtime identity tags.
// ReservedContextFieldPrefix is reserved for Durable Task context metadata.
const ReservedContextFieldPrefix = "__durabletask.context."

// ContextFields are immutable caller-supplied values propagated into task contexts.
type ContextFields map[string]string

// OrchestrationContextInfo identifies the orchestration associated with a task context.
// Orchestrator contexts contain the full persisted identity. Activity work items
// supply only InstanceID; the activity's own identity is in ActivityContextInfo.
type OrchestrationContextInfo struct {
InstanceID InstanceID
Name string
Expand Down
40 changes: 3 additions & 37 deletions internal/contextprop/tags.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,43 +7,9 @@ import (
"github.com/microsoft/durabletask-go/internal/tagcodec"
)

const (
instanceIDTag = api.ReservedContextFieldPrefix + "instance_id"
nameTag = api.ReservedContextFieldPrefix + "orchestration_name"
versionTag = api.ReservedContextFieldPrefix + "orchestration_version"
parentInstanceIDTag = api.ReservedContextFieldPrefix + "parent_instance_id"
)

// Encode returns a new tag map containing immutable fields and orchestration identity.
func Encode(
info api.OrchestrationContextInfo,
fields api.ContextFields,
userTags ...map[string]string,
) map[string]string {
tags := tagcodec.EncodeContextFields(fields)
if len(userTags) > 0 {
tags = tagcodec.Merge(tags, tagcodec.EncodeUserTags(userTags[0]))
}
if tags == nil {
tags = make(map[string]string, 5)
}
tags[tagcodec.ContextEncodingTag] = "1"
tags[instanceIDTag] = string(info.InstanceID)
tags[nameTag] = info.Name
tags[versionTag] = info.Version
tags[parentInstanceIDTag] = string(info.ParentInstanceID)
return tags
}

// Decode separates orchestration identity from caller-supplied immutable fields.
func Decode(tags map[string]string) (api.OrchestrationContextInfo, api.ContextFields) {
info := api.OrchestrationContextInfo{
InstanceID: api.InstanceID(tags[instanceIDTag]),
Name: tags[nameTag],
Version: tags[versionTag],
ParentInstanceID: api.InstanceID(tags[parentInstanceIDTag]),
}
return info, api.ContextFields(tagcodec.DecodeContextFields(tags))
// Encode returns a new tag map containing caller fields and user tags.
func Encode(fields api.ContextFields, userTags map[string]string) map[string]string {
return tagcodec.Merge(tagcodec.EncodeContextFields(fields), tagcodec.EncodeUserTags(userTags))
}

// Clone returns a defensive copy of tags, or nil when there is nothing to copy.
Expand Down
105 changes: 57 additions & 48 deletions internal/contextprop/tags_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,57 +5,66 @@ import (

"github.com/microsoft/durabletask-go/api"
"github.com/microsoft/durabletask-go/internal/tagcodec"
"github.com/stretchr/testify/require"
)

func TestEncodeDecode(t *testing.T) {
tags := Encode(api.OrchestrationContextInfo{
InstanceID: "instance",
Name: "orchestration",
Version: "v1",
ParentInstanceID: "parent",
}, api.ContextFields{"tenant": "alpha"})

info, fields := Decode(tags)
if info.InstanceID != "instance" ||
info.Name != "orchestration" ||
info.Version != "v1" ||
info.ParentInstanceID != "parent" {
t.Fatalf("unexpected info: %+v", info)
}

if fields["tenant"] != "alpha" {
t.Fatalf("tenant = %q, want alpha", fields["tenant"])
}
}

func TestEncodeSeparatesContextFieldsAndUserTags(t *testing.T) {
tags := Encode(
api.OrchestrationContextInfo{},
api.ContextFields{"tenant": "context"},
map[string]string{"team": "tag"},
)
_, fields := Decode(tags)
if fields["tenant"] != "context" {
t.Fatalf("tenant = %q, want context", fields["tenant"])
}
userTags := tagcodec.DecodeUserTags(tags)
if userTags["team"] != "tag" {
t.Fatalf("team = %q, want tag", userTags["team"])
}
if _, ok := fields["team"]; ok {
t.Fatalf("user tag leaked into context fields: %v", fields)
func TestEncodeCallerData(t *testing.T) {
for _, test := range []struct {
name string
fields api.ContextFields
tags map[string]string
want map[string]string
}{
{name: "absent"},
{name: "empty", fields: api.ContextFields{}, tags: map[string]string{}},
{
name: "user tags only",
tags: map[string]string{"tenant": ""},
want: map[string]string{tagcodec.ContextEncodingTag: "1", "tenant": ""},
},
{
name: "fields only",
fields: api.ContextFields{"tenant": ""},
want: map[string]string{tagcodec.ContextEncodingTag: "1", tagcodec.ContextFieldPrefix + "tenant": ""},
},
{
name: "same key in separate namespaces",
fields: api.ContextFields{"tenant": "context"},
tags: map[string]string{"tenant": "user"},
want: map[string]string{
tagcodec.ContextEncodingTag: "1",
tagcodec.ContextFieldPrefix + "tenant": "context",
"tenant": "user",
},
},
} {
t.Run(test.name, func(t *testing.T) {
encoded := Encode(test.fields, test.tags)
require.Equal(t, test.want, encoded)
fields := api.ContextFields(tagcodec.DecodeContextFields(encoded))
if len(test.fields) == 0 {
require.Nil(t, fields, "user tags must not become context fields")
} else {
require.Equal(t, test.fields, fields)
fields["tenant"] = "decoded mutation"
require.Equal(t, test.want, encoded)
}
if test.fields != nil {
test.fields["tenant"] = "caller mutation"
}
if test.tags != nil {
test.tags["tenant"] = "caller mutation"
}
require.Equal(t, test.want, encoded)
})
}
}

func TestEncodeOverwritesReservedCallerFields(t *testing.T) {
tags := Encode(api.OrchestrationContextInfo{}, api.ContextFields{
api.ReservedContextFieldPrefix + "orchestration_version": "spoofed",
})
info, fields := Decode(tags)
if info.Version != "" {
t.Fatalf("version = %q, want empty", info.Version)
}
if fields != nil {
t.Fatalf("reserved field leaked into caller fields: %v", fields)
}
func TestCloneDoesNotAliasTags(t *testing.T) {
require.Nil(t, Clone[map[string]string](nil))
require.Nil(t, Clone(map[string]string{}))
original := map[string]string{"empty": ""}
cloned := Clone(original)
cloned["empty"] = "changed"
require.Equal(t, map[string]string{"empty": ""}, original)
}
15 changes: 4 additions & 11 deletions task/context_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ func TestOrchestrationContextPropagatesOnlyPersistedIdentityAndFields(t *testing
nil,
wrapperspb.String("v2"),
)
started.GetExecutionStarted().Tags = contextprop.Encode(api.OrchestrationContextInfo{}, fields)
started.GetExecutionStarted().Tags = contextprop.Encode(fields, nil)
fields["tenant"] = "mutated"
events := []*protos.HistoryEvent{helpers.NewOrchestratorStartedEvent(), started}

Expand Down Expand Up @@ -164,7 +164,7 @@ func TestActivityContextPropagatesIdentityFieldsAndLogger(t *testing.T) {
}
}

func TestActivityContextDecodesDurableContextTags(t *testing.T) {
func TestActivityContextDecodesDurableContextFields(t *testing.T) {
registry := NewTaskRegistry()
if err := registry.AddActivityN("inspect-tags", func(ctx ActivityContext) (any, error) {
orchestration, _ := api.OrchestrationContextInfoFromContext(ctx.Context())
Expand All @@ -180,12 +180,7 @@ func TestActivityContextDecodesDurableContextTags(t *testing.T) {
}

event := helpers.NewTaskScheduledEvent(4, "inspect-tags", nil, nil, nil)
event.GetTaskScheduled().Tags = contextprop.Encode(api.OrchestrationContextInfo{
InstanceID: "tagged-instance",
Name: "tagged-parent",
Version: "v4",
ParentInstanceID: "root",
}, api.ContextFields{"tenant": "tagged"})
event.GetTaskScheduled().Tags = contextprop.Encode(api.ContextFields{"tenant": "tagged"}, nil)
response, err := NewTaskExecutor(registry).ExecuteActivity(
context.Background(),
"tagged-instance",
Expand All @@ -202,9 +197,7 @@ func TestActivityContextDecodesDurableContextTags(t *testing.T) {
if err := json.Unmarshal([]byte(response.GetTaskCompleted().GetResult().GetValue()), &output); err != nil {
t.Fatal(err)
}
if output.Orchestration.Name != "tagged-parent" ||
output.Orchestration.Version != "v4" ||
output.Orchestration.ParentInstanceID != "root" {
if output.Orchestration != (api.OrchestrationContextInfo{InstanceID: "tagged-instance"}) {
t.Fatalf("unexpected orchestration identity: %+v", output.Orchestration)
}
if output.Fields["tenant"] != "tagged" {
Expand Down
19 changes: 3 additions & 16 deletions task/executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,10 @@ import (
"time"

"github.com/microsoft/durabletask-go/api"
"github.com/microsoft/durabletask-go/internal/contextprop"
"github.com/microsoft/durabletask-go/internal/failure"
"github.com/microsoft/durabletask-go/internal/helpers"
"github.com/microsoft/durabletask-go/internal/protos"
"github.com/microsoft/durabletask-go/internal/tagcodec"
"google.golang.org/protobuf/types/known/timestamppb"
"google.golang.org/protobuf/types/known/wrapperspb"
)
Expand Down Expand Up @@ -187,23 +187,10 @@ func (te *taskExecutor) ExecuteActivity(ctx context.Context, id api.InstanceID,
}
ctx = api.ContextWithFields(ctx, te.contextFields)
ctx = helpers.ContextWithTraceContext(ctx, ts.GetParentTraceContext())
tagInfo, tagFields := contextprop.Decode(ts.GetTags())
ctx = api.ContextWithFields(ctx, tagFields)
ctx = api.ContextWithFields(ctx, api.ContextFields(tagcodec.DecodeContextFields(ts.GetTags())))
Comment thread
torosent marked this conversation as resolved.
orchestrationInfo, _ := api.OrchestrationContextInfoFromContext(ctx)
if orchestrationInfo.Name == "" {
orchestrationInfo.Name = tagInfo.Name
}
if orchestrationInfo.Version == "" {
orchestrationInfo.Version = tagInfo.Version
}
if orchestrationInfo.ParentInstanceID == "" {
orchestrationInfo.ParentInstanceID = tagInfo.ParentInstanceID
}
if orchestrationInfo.InstanceID == "" {
orchestrationInfo.InstanceID = tagInfo.InstanceID
if orchestrationInfo.InstanceID == "" {
orchestrationInfo.InstanceID = id
}
orchestrationInfo.InstanceID = id
}
ctx = api.WithOrchestrationContextInfo(ctx, orchestrationInfo)
ctx = api.WithActivityContextInfo(ctx, api.ActivityContextInfo{
Expand Down
27 changes: 4 additions & 23 deletions task/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -721,12 +721,8 @@ func (ctx *OrchestrationContext) internalScheduleActivity(
helpers.GetTaskFunctionName(activity),
options.rawInput,
options.versionOrInherited(ctx.Version))
scheduleTaskAction.GetScheduleTask().Tags = contextprop.Encode(api.OrchestrationContextInfo{
InstanceID: ctx.ID,
Name: ctx.Name,
Version: ctx.Version,
ParentInstanceID: ctx.parentInstanceID,
}, ctx.contextFields, mergeStringMaps(ctx.orchestrationTags, options.tags))
scheduleTaskAction.GetScheduleTask().Tags = contextprop.Encode(
ctx.contextFields, mergeStringMaps(ctx.orchestrationTags, options.tags))
Comment thread
torosent marked this conversation as resolved.

ctx.pendingActions[scheduleTaskAction.Id] = scheduleTaskAction

Expand Down Expand Up @@ -783,12 +779,6 @@ func (ctx *OrchestrationContext) internalCallSubOrchestrator(
options.versionOrDefault(ctx.defaultVersion),
)
createSubOrchestrationAction.GetCreateSubOrchestration().Tags = contextprop.Encode(
api.OrchestrationContextInfo{
InstanceID: ctx.ID,
Name: ctx.Name,
Version: ctx.Version,
ParentInstanceID: ctx.parentInstanceID,
},
mergeStringMaps(ctx.contextFields, options.contextFields),
mergeStringMaps(ctx.orchestrationTags, options.tags),
)
Expand Down Expand Up @@ -1277,7 +1267,7 @@ func (ctx *OrchestrationContext) onExecutionStarted(es *protos.ExecutionStartedE
ctx.Name = es.Name
ctx.Version = es.GetVersion().GetValue()
ctx.executionID = es.GetOrchestrationInstance().GetExecutionId().GetValue()
_, fields := contextprop.Decode(es.GetTags())
fields := api.ContextFields(tagcodec.DecodeContextFields(es.GetTags()))
ctx.contextFields = mergeStringMaps(ctx.contextFields, fields)
ctx.orchestrationTags = tagcodec.DecodeUserTagsOrPlain(es.GetTags())
if parent := es.GetParentInstance(); parent != nil {
Expand Down Expand Up @@ -1779,16 +1769,7 @@ func (ctx *OrchestrationContext) setCompleteInternal(
completed := completedAction.GetCompleteOrchestration()
if status == protos.OrchestrationStatus_ORCHESTRATION_STATUS_CONTINUED_AS_NEW {
completed.NewVersion = ctx.continuedAsNewVersion
version := ctx.Version
if ctx.continuedAsNewVersion != nil {
version = ctx.continuedAsNewVersion.GetValue()
}
completed.Tags = contextprop.Encode(api.OrchestrationContextInfo{
InstanceID: ctx.ID,
Name: ctx.Name,
Version: version,
ParentInstanceID: ctx.parentInstanceID,
}, ctx.contextFields, ctx.orchestrationTags)
completed.Tags = contextprop.Encode(ctx.contextFields, ctx.orchestrationTags)
} else {
completed.Tags = tagcodec.EncodeUserTags(ctx.orchestrationTags)
}
Expand Down
Loading
Loading