Add stateful-history delta work items to the workflow worker - #1777
Add stateful-history delta work items to the workflow worker#1777JoshVanL wants to merge 3 commits into
Conversation
The sidecar re-sends a workflow instance's entire committed history to the worker on every turn. This adds the worker half of the "stateful history" optimization so that, once a worker is warm for an instance on a work-item stream, the sidecar sends only the new committed events (the delta) and the worker reconstructs the full history from its own cache. It mirrors the Go (durabletask-go), Python, and .NET SDK implementations and is on by default. Worker (durabletask-client): - WorkflowHistoryCache: a per-stream cache of each instance's committed history, bounded by a sliding TTL, an instance-count cap, and a byte budget with LRU eviction. Injectable clock for deterministic tests. - DurableTaskGrpcWorker: advertise WORKER_CAPABILITY_STATEFUL_HISTORY in GetWorkItemsRequest, reset the cache on every reconnect (the sidecar drops the old stream's warm set), and reclaim idle entries with a daemon janitor stopped on close. - OrchestratorRunner: before replay, resolve the full committed history (cached prefix + delta on a hit, or a GetInstanceHistory fetch on a miss) instead of using the request's pastEvents directly; after replay, cache the committed history, or drop it once the instance ends (a CompleteWorkflow action, covering completed/failed/terminated/continued-as-new). A TerminateWorkflow action targets a different instance and is deliberately not treated as a reset. Correctness never depends on the cache: any miss (cold stream, eviction, desync) self-heals via the GetInstanceHistory fallback, so this only changes per-turn bandwidth, not results. A fallback fetch that fails abandons the work item for backend redelivery rather than completing with a partial history. Configuration (DurableTaskGrpcWorkerBuilder): - disableStatefulHistory to opt out, plus historyCacheTtl, historyCacheMaxInstances, and historyCacheMaxBytes to tune the bounds. Signed-off-by: joshvanl <me@joshvanl.dev>
7a6842a to
2ac527f
Compare
There was a problem hiding this comment.
Pull request overview
Adds the worker-side implementation of the “stateful history” optimization so the sidecar can send only committed-history deltas on a warm work-item stream, with safe self-healing via GetInstanceHistory on cache misses.
Changes:
- Introduces a per-stream
WorkflowHistoryCache(TTL + LRU bounds) and uses it to reconstruct full committed history from deltas. - Updates
DurableTaskGrpcWorker/OrchestratorRunnerto advertise the capability, reset cache on reconnect, and abandon (drop stream) when history recovery fails. - Adds unit + worker-level + end-to-end integration tests validating cache behavior and actual on-the-wire delta delivery.
Reviewed changes
Copilot reviewed 12 out of 12 changed files in this pull request and generated 3 comments.
Show a summary per file
| File | Description |
|---|---|
| sdk-workflows/src/main/java/io/dapr/workflows/runtime/WorkflowRuntimeBuilder.java | Exposes stateful-history configuration knobs on the workflow runtime builder and forwards them to the worker builder. |
| sdk-workflows/src/test/java/io/dapr/workflows/runtime/WorkflowRuntimeBuilderTest.java | Verifies default enablement and forwarding of stateful-history options (via reflection). |
| durabletask-client/src/main/java/io/dapr/durabletask/DurableTaskGrpcWorkerBuilder.java | Adds stateful-history configuration fields and builder methods. |
| durabletask-client/src/main/java/io/dapr/durabletask/DurableTaskGrpcWorker.java | Advertises capability, manages per-stream cache lifecycle, and runs a janitor sweep for TTL eviction. |
| durabletask-client/src/main/java/io/dapr/durabletask/WorkflowHistoryCache.java | Implements the per-stream committed-history cache with TTL/LRU bounds and byte accounting. |
| durabletask-client/src/main/java/io/dapr/durabletask/runner/OrchestratorRunner.java | Resolves full committed history using cached prefix + delta or fallback fetch, and updates/evicts cache entries after turns. |
| durabletask-client/src/test/java/io/dapr/durabletask/WorkItemObserver.java | Adds a gRPC interceptor to count full-sends vs deltas and history fetches for wire-level assertions. |
| durabletask-client/src/test/java/io/dapr/durabletask/WorkflowHistoryCacheTest.java | Unit tests for cache eviction policies, TTL sliding behavior, and immutability/snapshotting behavior. |
| durabletask-client/src/test/java/io/dapr/durabletask/StatefulHistoryIT.java | Integration test validating that the sidecar actually sends deltas on the wire and that warm streams avoid cache-miss fetches. |
| durabletask-client/src/test/java/io/dapr/durabletask/runner/OrchestratorRunnerHistoryTest.java | Deterministic tests for history resolution and cache update behavior in the runner. |
| durabletask-client/src/test/java/io/dapr/durabletask/IntegrationTestBase.java | Extends test worker builder to accept a custom gRPC channel and toggle stateful history for integration tests. |
| durabletask-client/src/test/java/io/dapr/durabletask/DurableTaskGrpcWorkerStatefulHistoryTest.java | Worker-level tests against an in-process fake sidecar covering capability advertisement and miss-recovery behavior. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
Signed-off-by: joshvanl <me@joshvanl.dev>
Signed-off-by: joshvanl <me@joshvanl.dev>
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #1777 +/- ##
============================================
+ Coverage 77.10% 77.13% +0.02%
- Complexity 2317 2321 +4
============================================
Files 245 245
Lines 7186 7194 +8
Branches 750 750
============================================
+ Hits 5541 5549 +8
Misses 1284 1284
Partials 361 361 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
| void updateHistoryCache( | ||
| String instanceId, | ||
| List<HistoryEvents.HistoryEvent> pastEvents, | ||
| TaskOrchestratorResult result) { | ||
| if (this.historyCache == null) { | ||
| return; | ||
| } | ||
|
|
||
| boolean ended = result.getActions().stream() | ||
| .anyMatch(OrchestratorActions.WorkflowAction::hasCompleteWorkflow); | ||
| if (ended) { | ||
| this.historyCache.remove(instanceId); | ||
| } else { | ||
| this.historyCache.put(instanceId, pastEvents); | ||
| } | ||
| } |
There was a problem hiding this comment.
It seems we're not porting one check from the Go reference. We write to the cache
unconditionally, while Go first checks that the turn doing the writing is still the
current one (worker_grpc.go:245).
Without it, a runner that stalls across a reconnect can overwrite a newer entry.
Most of the time the length check catches it and we just refetch, but after a
continue-as-new the stale prefix can have the same length as the good one, and
then it goes through silently.
Possible fix: port latestTokens / noteDispatch / isLatestDispatch /
forgetDispatch from worker_history.go:120-186, call noteDispatch where the
work item is dispatched, and gate the updateHistoryCache body on it. The one
thing not to get wrong: the token map must not be cleared in reset().
| synchronized (this.lock) { | ||
| Entry existing = this.entries.get(instanceId); | ||
| if (existing != null) { | ||
| this.totalBytes -= existing.bytes; | ||
| } | ||
| this.entries.put(instanceId, new Entry(snapshot, bytes, this.clockNanos.getAsLong())); | ||
| this.totalBytes += bytes; | ||
| this.evictToFit(instanceId); | ||
| } | ||
| } |
There was a problem hiding this comment.
Should be check if bytes are 0 so we do not put it on the entries?
The sidecar re-sends a workflow instance's entire committed history to the worker on every turn. This adds the worker half of the "stateful history" optimization so that, once a worker is warm for an instance on a work-item stream, the sidecar sends only the new committed events (the delta) and the worker reconstructs the full history from its own cache. It mirrors the Go (durabletask-go), Python, and .NET SDK implementations and is on by default.
Worker (durabletask-client):
Correctness never depends on the cache: any miss (cold stream, eviction, desync) self-heals via the GetInstanceHistory fallback, so this only changes per-turn bandwidth, not results. A fallback fetch that fails abandons the work item for backend redelivery rather than completing with a partial history.
Configuration (DurableTaskGrpcWorkerBuilder):