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

### Fixed

- Forced coroutine unloading no longer leaks durable actions or custom status from application defers: durable operations panic with `ErrTaskBlocked` before changing state. Normal-return defers remain awaitable across replay, and orchestration loggers suppress forced-unload output.
- Long durable timers no longer schedule another chunk when a trailing timer event is processed after the orchestration has finalized.
- DTS authentication now caches access tokens and immutable gRPC metadata, honors credential `RefreshOn` guidance, and coalesces concurrent refreshes and failures. Previously credentials such as `AzureCLICredential` were invoked for every RPC, serializing high-throughput workers behind external token acquisition.
- Disabled large-payload handling no longer builds transform closures, maps, and goroutines for ordinary payloads. Reserved payload-reference prefixes are still rejected when no store and resolver are configured.
- Worker completion and abandon RPCs no longer retry `NotFound` responses ten times. A completion whose work item is already unavailable is dropped for normal DTS redelivery without a second, guaranteed-failing abandon RPC.
Expand Down
57 changes: 57 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -307,6 +307,63 @@ Obey these rules when you write a retry handler:

Full sample: [samples/retries](./samples/retries).

### Deferred work and cooperative cleanup

A root or child coroutine's defers can schedule and await durable work during
normal return, error return, or an application panic. A deferred `Await` can
span turns: replay resumes the cleanup before that coroutine completes.

Forced unloading is different. At the end of a waiting turn, or when the
orchestration ends with waiting children, Go unwinds their stacks. Durable
operations in those defers panic with `task.ErrTaskBlocked` before changing
durable state. Never recover this signal. The rest of that deferred function
is skipped, but other defers still run. Termination does not run business cleanup.

`ctx.Logger()` suppresses replay and forced-unload output without changing
`ctx.IsReplaying`. Other local defers, such as metrics or `span.End`, can run
every turn. Keep local-resource cleanup separate from durable cleanup.

To finish child cleanup before completion, run the child on the uncanceled
root, cancel only its body operations, and join it. A coroutine canceled before
it starts is skipped, so do not rely on a `Done` deferred inside `child.Go`.

```go
func CleanupOrchestrator(ctx *task.OrchestrationContext) (any, error) {
work, cancel := ctx.WithCancel()
group := ctx.NewWaitGroup()
var childErr error
group.Add(1)
ctx.Go(func(cleanup *task.OrchestrationContext) {
defer group.Done()
defer func() {
childErr = errors.Join(childErr, cleanup.CallActivity("Cleanup").Await(nil))
}()
childErr = work.WaitForSingleEvent("work", -1).Await(nil)
if childErr == task.ErrTaskCanceled {
childErr = nil // The body cancellation is expected; cleanup errors are not.
}
})
if err := ctx.WaitForSingleEvent("stop", -1).Await(nil); err != nil {
return nil, err
}
cancel()
group.Wait(ctx)
return nil, childErr
}
```

The `Cleanup` activity must act only on resources this workflow owns and must
be idempotent, since activity delivery is at least once. For multiple children,
keep separate error slots and combine them after joining.

Unlike `Await`, `Select` and `WaitGroup.Wait` panic on cancellation; wrap those
body operations with a handler that recovers only the exact
`task.ErrTaskCanceled` value and re-panics everything else.

Use a deferred function when scheduling itself should be deferred.
`defer ctx.CallActivity("Cleanup").Await(nil)` schedules the activity immediately
and defers only the await.

### Orchestration management

Use a `TaskRegistry` to register your orchestrator, activity, and entity functions. Then use the client from `durabletaskscheduler.NewClient` to control the orchestrations.
Expand Down
5 changes: 3 additions & 2 deletions task/context.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,8 @@ func (ctx *OrchestrationContext) Context() context.Context {
})
}

// Logger returns a slog logger that suppresses output while replaying history.
// Logger returns a slog logger that suppresses output while replaying history
// and while waiting coroutines are forcibly unloaded.
func (ctx *OrchestrationContext) Logger() *slog.Logger {
engine := ctx.engineContext()
logger := engine.logger
Expand All @@ -46,7 +47,7 @@ func (ctx *OrchestrationContext) Logger() *slog.Logger {
handler := &replaySafeHandler{
handler: logger.Handler(),
replaying: func() bool {
return engine.IsReplaying
return engine.IsReplaying || engine.unwinding
},
}
return slog.New(handler).With(
Expand Down
4 changes: 2 additions & 2 deletions task/eventchannel.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ func NewEventChannel[T any](ctx *OrchestrationContext, name string) *EventChanne
if ctx == nil {
panic("event channel requires an orchestration context")
}
engine := ctx.engineContext()
engine := ctx.effectContext()
key := strings.ToUpper(name)
if existing, ok := engine.eventChannels[key]; ok {
channel, ok := existing.(*EventChannel[T])
Expand Down Expand Up @@ -57,7 +57,7 @@ func (c *EventChannel[T]) Receive(ctx *OrchestrationContext) T {
// ReceiveErr waits for and consumes the next event value, returning payload
// decoding and cancellation errors to the orchestrator.
func (c *EventChannel[T]) ReceiveErr(ctx *OrchestrationContext) (T, error) {
if ctx.engineContext() != c.ctx {
if ctx.effectContext() != c.ctx {
panic("event channel used with a different orchestration context")
}
var value T
Expand Down
53 changes: 38 additions & 15 deletions task/orchestrator.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,11 @@ type replayEvent struct {
}

// OrchestrationContext is the parameter type for orchestrator functions.
//
// Deferred functions can issue and await durable work when their coroutine
// returns normally. When the runtime forcibly unloads a waiting coroutine,
// durable operations instead panic with ErrTaskBlocked before changing state.
// Applications must not recover that control-flow signal.
type OrchestrationContext struct {
ID api.InstanceID
Name string
Expand Down Expand Up @@ -75,6 +80,7 @@ type OrchestrationContext struct {
converter api.DataConverter
entitiesSupported bool
scheduler *coroutineScheduler
unwinding bool
root *OrchestrationContext
scope *cancellationScope
derived []*OrchestrationContext
Expand Down Expand Up @@ -123,13 +129,15 @@ type ContinueAsNewOption func(*OrchestrationContext)
// runtime to carry forward any unprocessed external events to the new instance.
func WithKeepUnprocessedEvents() ContinueAsNewOption {
return func(ctx *OrchestrationContext) {
ctx.effectContext()
ctx.saveBufferedExternalEvents = true
}
}

// WithContinueAsNewVersion migrates the next execution to a new orchestration version.
func WithContinueAsNewVersion(version string) ContinueAsNewOption {
return func(ctx *OrchestrationContext) {
ctx.effectContext()
if version != "" && strings.TrimSpace(version) == "" {
ctx.continuedAsNewVersion = nil
return
Expand Down Expand Up @@ -289,6 +297,15 @@ func (ctx *OrchestrationContext) engineContext() *OrchestrationContext {
return ctx
}

// effectContext stops forced-unwind defers before they can mutate durable state.
func (ctx *OrchestrationContext) effectContext() *OrchestrationContext {
engine := ctx.engineContext()
if engine.unwinding {
panic(ErrTaskBlocked)
}
return engine
}

func (ctx *OrchestrationContext) syncDerivedContexts() {
active := ctx.derived[:0]
for _, derived := range ctx.derived {
Expand All @@ -310,8 +327,10 @@ func (ctx *OrchestrationContext) syncDerivedContexts() {
// WithCancel creates a child orchestration context whose tasks, nested scopes,
// and coroutines are canceled together at the next scheduler step.
// A child of an already-canceled scope is canceled immediately.
// Await cleanup from an uncanceled wrapper, cancel only its body operations,
// and join the wrapper before returning.
func (ctx *OrchestrationContext) WithCancel() (*OrchestrationContext, func()) {
engine := ctx.engineContext()
engine := ctx.effectContext()
if engine.scheduler == nil {
panic("cancellation scope created outside orchestrator execution")
}
Expand Down Expand Up @@ -583,12 +602,12 @@ func (ctx *OrchestrationContext) processEvent(e *protos.HistoryEvent) error {
// SetCustomStatus stores a raw, pre-serialized custom status string.
// Use SetCustomStatusValue to apply the configured data converter.
func (octx *OrchestrationContext) SetCustomStatus(cs string) {
octx.engineContext().customStatus = cs
octx.effectContext().customStatus = cs
}

// SetCustomStatusValue serializes and stores a typed custom status value.
func (octx *OrchestrationContext) SetCustomStatusValue(value any) error {
engine := octx.engineContext()
engine := octx.effectContext()
payload, err := api.SerializeData(engine.converter, value)
if err != nil {
return fmt.Errorf("failed to serialize custom status: %w", err)
Expand All @@ -599,14 +618,14 @@ func (octx *OrchestrationContext) SetCustomStatusValue(value any) error {

// SetRawCustomStatus stores a pre-serialized custom status value.
func (octx *OrchestrationContext) SetRawCustomStatus(payload string) {
octx.engineContext().customStatus = payload
octx.effectContext().customStatus = payload
}

var guidNamespace = uuid.MustParse("9e952958-5e33-4daf-827f-2fa12937b875")

// NewGuid returns a deterministic UUID that is stable across orchestration replay.
func (ctx *OrchestrationContext) NewGuid() string {
engine := ctx.engineContext()
engine := ctx.effectContext()
timestamp := engine.CurrentTimeUtc.UTC().Format("2006-01-02T15:04:05.0000000Z")
name := fmt.Sprintf("%s_%s_%d", engine.ID, timestamp, engine.newGuidCounter)
engine.newGuidCounter++
Expand All @@ -623,7 +642,7 @@ func (octx *OrchestrationContext) GetInput(v any) error {
// parameter can be either the name of an activity as a string or can be a pointer to the function
// that implements the activity, in which case the name is obtained via reflection.
func (ctx *OrchestrationContext) CallActivity(activity any, opts ...CallActivityOption) Task {
engine := ctx.engineContext()
engine := ctx.effectContext()
options := new(callActivityOptions)
for _, configure := range opts {
if err := configure(options, engine.converter); err != nil {
Expand Down Expand Up @@ -659,6 +678,7 @@ func (ctx *OrchestrationContext) newFailedTask(engine *OrchestrationContext, err
// Go starts a coroutine that is cooperatively scheduled with the orchestration.
// Only one orchestration coroutine runs at a time, in monotonically increasing ID order.
// A callback whose scope is canceled before it starts is not invoked.
// Join children explicitly when their cleanup must finish before the root returns.
func (ctx *OrchestrationContext) Go(fn func(ctx *OrchestrationContext)) {
if fn == nil {
panic("orchestration coroutine function must be non-nil")
Expand Down Expand Up @@ -712,7 +732,7 @@ func (ctx *OrchestrationContext) internalScheduleActivity(
}

func (ctx *OrchestrationContext) CallSubOrchestrator(orchestrator any, opts ...SubOrchestratorOption) Task {
engine := ctx.engineContext()
engine := ctx.effectContext()
if engine.criticalSectionID != "" {
return ctx.newFailedTask(engine, fmt.Errorf("sub-orchestrations cannot be started while holding entity locks"))
}
Expand Down Expand Up @@ -902,7 +922,7 @@ func computeNextDelay(currentTimeUtc time.Time, policy RetryPolicy, attempt int,

// CreateTimer schedules a durable timer that expires after the specified delay.
func (ctx *OrchestrationContext) CreateTimer(delay time.Duration) Task {
engine := ctx.engineContext()
engine := ctx.effectContext()
if ctx.scope.isCanceled() {
return newTaskInScope(engine, ctx.scope)
}
Expand All @@ -926,7 +946,7 @@ func (ctx *OrchestrationContext) createTimerInternal(

var scheduleNextChunk func()
scheduleNextChunk = func() {
if logicalTimer.isCompleted {
if logicalTimer.isCompleted || (ctx.scheduler != nil && ctx.scheduler.isStopping()) {
return
}

Expand Down Expand Up @@ -988,7 +1008,7 @@ func (ctx *OrchestrationContext) createTimerAction(
//
// Note that event names are case-insensitive.
func (ctx *OrchestrationContext) WaitForSingleEvent(eventName string, timeout time.Duration) Task {
engine := ctx.engineContext()
engine := ctx.effectContext()
task := newTaskInScope(engine, ctx.scope)
if ctx.scope.isCanceled() {
return task
Expand Down Expand Up @@ -1023,7 +1043,7 @@ func (ctx *OrchestrationContext) WaitForSingleEvent(eventName string, timeout ti

// CallEntity sends an operation request to an entity and waits for its response.
func (ctx *OrchestrationContext) CallEntity(entityID api.EntityID, operationName string, opts ...callEntityOption) Task {
engine := ctx.engineContext()
engine := ctx.effectContext()
if engine.isTerminated || ctx.scope.isCanceled() {
task := newTaskInScope(engine, ctx.scope)
task.cancel()
Expand Down Expand Up @@ -1086,7 +1106,7 @@ func (ctx *OrchestrationContext) CallEntity(entityID api.EntityID, operationName

// SignalEntity sends a fire-and-forget entity operation.
func (ctx *OrchestrationContext) SignalEntity(entityID api.EntityID, operationName string, opts ...signalEntityOption) error {
engine := ctx.engineContext()
engine := ctx.effectContext()
if engine.isTerminated || ctx.scope.isCanceled() {
return ErrTaskCanceled
}
Expand Down Expand Up @@ -1128,7 +1148,7 @@ func (ctx *OrchestrationContext) SignalEntity(entityID api.EntityID, operationNa
// If cancellation follows a committed request, critical-section restrictions
// remain in effect until the eventual grant is received and automatically released.
func (ctx *OrchestrationContext) LockEntities(entityIDs ...api.EntityID) (func(), error) {
engine := ctx.engineContext()
engine := ctx.effectContext()
if engine.isTerminated || ctx.scope.isCanceled() {
return nil, ErrTaskCanceled
}
Expand Down Expand Up @@ -1211,8 +1231,9 @@ func (ctx *OrchestrationContext) IsInCriticalSection() bool {
return ctx.engineContext().criticalSectionID != ""
}

// ContinueAsNew requests a new execution when the root coroutine finishes.
func (ctx *OrchestrationContext) ContinueAsNew(newInput any, options ...ContinueAsNewOption) {
engine := ctx.engineContext()
engine := ctx.effectContext()
engine.continuedAsNew = true
engine.continuedAsNewInput = newInput
for _, option := range options {
Expand All @@ -1223,7 +1244,7 @@ func (ctx *OrchestrationContext) ContinueAsNew(newInput any, options ...Continue
// SendEvent sends an event to another orchestration instance as part of the
// current durable orchestration transaction.
func (ctx *OrchestrationContext) SendEvent(instanceID api.InstanceID, eventName string, payload any) error {
engine := ctx.engineContext()
engine := ctx.effectContext()
raw, err := marshalData(engine.converter, payload)
if err != nil {
return fmt.Errorf("failed to marshal event payload: %w", err)
Expand Down Expand Up @@ -1426,6 +1447,7 @@ func (ctx *OrchestrationContext) peekBufferedEvent(key string) (*bufferedEvent,
}

func (ctx *OrchestrationContext) takeBufferedEvent(key string) (*bufferedEvent, bool) {
ctx.effectContext()
eventList, ok := ctx.bufferedExternalEvents[key]
if !ok || eventList.Len() == 0 {
return nil, false
Expand Down Expand Up @@ -1788,6 +1810,7 @@ func (ctx *OrchestrationContext) clearCriticalSection() {
}

func (ctx *OrchestrationContext) getNextSequenceNumber() int32 {
ctx.effectContext()
current := ctx.sequenceNumber
ctx.sequenceNumber++
return current
Expand Down
4 changes: 4 additions & 0 deletions task/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,10 @@ func (s *coroutineScheduler) shutdown() {
return
}
s.stopping = true
engine := s.ctx.engineContext()
// Stopping persists while trailing history is processed; unwinding must not.
engine.unwinding = true
defer func() { engine.unwinding = false }()
for _, c := range s.all {
c.exit()
}
Expand Down
2 changes: 1 addition & 1 deletion task/select.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,10 +94,10 @@ func (ctx *OrchestrationContext) Select(cases ...SelectCase) {
}

func (ctx *OrchestrationContext) selectCase(cases []SelectCase) SelectCase {
engine := ctx.effectContext()
if len(cases) == 0 {
panic("Select requires at least one case")
}
engine := ctx.engineContext()
scheduler := engine.scheduler
if scheduler == nil {
panic("Select called outside orchestrator execution")
Expand Down
3 changes: 2 additions & 1 deletion task/select_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -389,9 +389,10 @@ func executeOrchestrationTurn(
instanceID api.InstanceID,
oldEvents []*protos.HistoryEvent,
newEvents []*protos.HistoryEvent,
options ...TaskExecutorOption,
) *protos.OrchestratorResponse {
t.Helper()
result, err := NewTaskExecutor(registry).ExecuteOrchestrator(
result, err := NewTaskExecutor(registry, options...).ExecuteOrchestrator(
context.Background(),
instanceID,
oldEvents,
Expand Down
8 changes: 6 additions & 2 deletions task/task.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ import (
// ErrTaskBlocked is not an error, but rather a control flow signal indicating that an orchestrator
// function has executed as far as it can and that it now needs to unload, dispatch any scheduled tasks,
// and commit its current execution progress to durable storage.
// Durable operations also panic with this signal when called by deferred code
// during forced coroutine unloading, before creating or modifying durable work.
var ErrTaskBlocked = errors.New("the current task is blocked")

// ErrTaskCanceled is used to indicate that a task was canceled. Tasks can be canceled, for example,
Expand Down Expand Up @@ -45,6 +47,7 @@ type completableTask struct {
}

func newTaskInScope(ctx *OrchestrationContext, scope *cancellationScope) *completableTask {
ctx.effectContext()
task := &completableTask{
orchestrationCtx: ctx,
waiters: make(map[*coroutine]struct{}),
Expand All @@ -62,11 +65,13 @@ func newTaskInScope(ctx *OrchestrationContext, scope *cancellationScope) *comple
//
// Await will return ErrTaskCanceled if the task was canceled - e.g. due to a timeout.
//
// Await may panic with ErrTaskBlocked as the panic value if called on a task that has not yet completed.
// Await may panic with ErrTaskBlocked as the panic value if called on a task that has not yet completed,
// or if called while the runtime forcibly unloads the current coroutine.
// This is normal control flow behavior for orchestrator functions and doesn't actually indicate a failure
// of any kind. However, orchestrator functions must never attempt to recover from such panics to ensure that
// the orchestration execution can procede normally.
func (t *completableTask) Await(v any) error {
t.orchestrationCtx.effectContext()
for {
if t.isCompleted {
if t.localErr != nil {
Expand Down Expand Up @@ -116,7 +121,6 @@ func (t *completableTask) Await(v any) error {
break
}
}
// TODO: Need a rule about using "defer" in orchestrations because planned panics will invoke them unexpectedly
panic(ErrTaskBlocked)
}

Expand Down
Loading
Loading