Skip to content

fix(temporal): send activity heartbeats from inside activities - #540

Open
michaelxu2288 wants to merge 2 commits into
scaleapi:mainfrom
michaelxu2288:fix/activity-heartbeats
Open

michaelxu2288 wants to merge 2 commits into
scaleapi:mainfrom
michaelxu2288:fix/activity-heartbeats

Conversation

@michaelxu2288

@michaelxu2288 michaelxu2288 commented Oct 1, 2026 •

Copy link
Copy Markdown

Problem

heartbeat_if_in_workflow() in src/agentex/lib/utils/temporal.py never sends a heartbeat:

def heartbeat_if_in_workflow(heartbeat_name: str):
    if in_temporal_workflow():          # workflow.in_workflow()
        activity.heartbeat(heartbeat_name)

activity.heartbeat() is only legal inside an activity, and inside an activity workflow.in_workflow() is False. 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_timeout equal to start_to_close_timeout. An agent that sets a shorter heartbeat_timeout on a long activity gets ActivityTaskTimedOut (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

  • Add heartbeat_if_in_activity(). It heartbeats when activity.in_activity() is true and is a silent no-op everywhere else, including workflow code and sync agents.
  • Keep 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, with on_heartbeat collecting heartbeats from an activity that calls the helper:
    • On main: heartbeat_if_in_workflow("doing slow work") records [].
    • On this branch: it records [('doing slow work',)]. The new heartbeat_if_in_activity records the same.
  • New tests in tests/lib/test_temporal_utils.py:
    • one heartbeat inside an activity, for both names;
    • a silent no-op outside an activity.
  • uv run pytest -n 0 tests/lib/test_temporal_utils.py tests/lib/core/services tests/lib/core/temporal: 74 passed.
  • ruff check and pyright are clean on both files.

RetriggerConfidence Score: 4/5

The PR is not ready to merge because a long pause in a streamed turn can still time out the activity.

Fix All in CursorFindings

  1. P1 Live OpenAI streams time out ▶
Fix with agent prompt
### Issue 1
src/agentex/lib/utils/temporal.py:22-23
`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.

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.

  • Activity heartbeats now reach Temporal from activity code.
  • Streamed-turn delivery sends heartbeats as events arrive.

Reviews (2) · Last reviewed commit: "fix(harness): heartbeat while auto_send ..."

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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 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.

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.

Fix in Cursor Fix in Claude Code Fix in Codex

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant