fix(temporal): send activity heartbeats from inside activities - #540
Open
michaelxu2288 wants to merge 2 commits into
Open
michaelxu2288 wants to merge 2 commits into
michaelxu2288 wants to merge 2 commits into
Conversation
heartbeat_if_in_workflow() only called activity.heartbeat() when workflow.in_workflow() was true. activity.heartbeat() is only legal inside an activity, where workflow.in_workflow() is false, so the guard never passed where it mattered. Every ADK service that calls it (tasks, messages, streaming, the LiteLLM/OpenAI/SGP providers, ACP, tracing, templating) runs inside an activity and never heartbeated. The shipped defaults hide this because each provider sets heartbeat_timeout equal to start_to_close_timeout. An agent that sets a shorter heartbeat_timeout for a long activity, the usual Temporal pattern, gets ActivityTaskTimedOut (HEARTBEAT) on every attempt, and activity cancellation is never delivered through the heartbeat channel. Add heartbeat_if_in_activity(), which heartbeats when activity.in_activity() is true and is a no-op everywhere else, and keep heartbeat_if_in_workflow() as an alias so the existing call sites work unchanged. Verified with temporalio.testing.ActivityEnvironment: inside an activity the old helper records no heartbeat, the new one (and the alias) records one. tests/lib/test_temporal_utils.py, tests/lib/core/services and tests/lib/core/temporal: 74 passed.
Comment on lines
+22
to
23
| if activity.in_activity(): | ||
| activity.heartbeat(heartbeat_name) |
There was a problem hiding this comment.
run_agent_streamed_auto_send now sends a heartbeat at the start, but sends none while it streams. If a caller sets heartbeat_timeout shorter than the stream’s duration, Temporal can time out and retry an activity that is still sending messages. This path needs heartbeats while the stream runs.
Prompt To Fix With AI
This is a comment left during a code review.
Path: src/agentex/lib/utils/temporal.py
Line: 22-23
Comment:
**Live OpenAI streams time out**
`run_agent_streamed_auto_send` now sends a heartbeat at the start, but sends none while it streams. If a caller sets `heartbeat_timeout` shorter than the stream’s duration, Temporal can time out and retry an activity that is still sending messages. This path needs heartbeats while the stream runs.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.
Author
There was a problem hiding this comment.
Fixed in 62b0b41, one level down so it covers every harness path: auto_send now heartbeats for each event it delivers (temporalio throttles the RPCs; no-op outside an activity). Test runs it in an ActivityEnvironment and records one heartbeat per event.
The provider helpers heartbeat once when a streamed turn starts, then hand the stream to UnifiedEmitter.auto_send_turn, which never heartbeated. A streamed OpenAI turn (or any harness turn delivered from an activity) that outlasts a short heartbeat_timeout was timed out and retried while it was still sending messages. auto_send now calls heartbeat_if_in_activity for each event it delivers; temporalio throttles the actual RPCs, and it is a no-op outside an activity. New test: auto_send inside a temporalio ActivityEnvironment records one heartbeat per event (none before this change).
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
heartbeat_if_in_workflow()insrc/agentex/lib/utils/temporal.pynever sends a heartbeat:activity.heartbeat()is only legal inside an activity, and inside an activityworkflow.in_workflow()isFalse. So the guard never passes where a heartbeat is possible. The helper is called 36 times across the ADK services: tasks, messages, streaming, the LiteLLM/OpenAI/SGP providers (including inside their streaming loops), ACP, tracing and templating. All of those run inside activities, and none of them heartbeat.The shipped defaults hide this because each provider module sets
heartbeat_timeoutequal tostart_to_close_timeout. An agent that sets a shorterheartbeat_timeouton a long activity getsActivityTaskTimedOut (HEARTBEAT)on every attempt, which is the standard Temporal pattern for long LLM calls. It also means activity cancellation is never delivered through the heartbeat channel.Fix
heartbeat_if_in_activity(). It heartbeats whenactivity.in_activity()is true and is a silent no-op everywhere else, including workflow code and sync agents.heartbeat_if_in_workflow()as an alias, so every existing call site starts heartbeating without being touched. Renaming the call sites can be a follow-up if you want the old name gone.One behaviour change to expect: with heartbeats flowing, activities now receive Temporal's cancellation when their workflow cancels them, instead of running to completion.
Verification
temporalio.testing.ActivityEnvironment, withon_heartbeatcollecting heartbeats from an activity that calls the helper:main:heartbeat_if_in_workflow("doing slow work")records[].[('doing slow work',)]. The newheartbeat_if_in_activityrecords the same.tests/lib/test_temporal_utils.py:uv run pytest -n 0 tests/lib/test_temporal_utils.py tests/lib/core/services tests/lib/core/temporal: 74 passed.ruff checkandpyrightare clean on both files.The PR is not ready to merge because a long pause in a streamed turn can still time out the activity.
Fix with agent prompt
Summary
The PR routes shared Temporal heartbeat calls through the activity context and adds a heartbeat for every event delivered by
auto_send(). The old helper name stays available, so existing callers use the corrected behavior.Reviews (2) · Last reviewed commit: "fix(harness): heartbeat while auto_send ..."