Skip to content

Commit 9ddbdd4

Browse files
committed
fix(temporal): send activity heartbeats from inside activities
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.
1 parent 6d68f3a commit 9ddbdd4

2 files changed

Lines changed: 62 additions & 2 deletions

File tree

‎src/agentex/lib/utils/temporal.py‎

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,11 +12,26 @@ def in_temporal_workflow():
1212
return False
1313

1414

15-
def heartbeat_if_in_workflow(heartbeat_name: str):
16-
if in_temporal_workflow():
15+
def heartbeat_if_in_activity(heartbeat_name: str) -> None:
16+
"""Records a Temporal activity heartbeat when called from inside an activity.
17+
18+
``activity.heartbeat`` is only legal inside an activity, so this stays a
19+
silent no-op everywhere else: workflow code and sync agents reach the same
20+
shared services and must not raise there.
21+
"""
22+
if activity.in_activity():
1723
activity.heartbeat(heartbeat_name)
1824

1925

26+
def heartbeat_if_in_workflow(heartbeat_name: str) -> None:
27+
"""Deprecated alias for :func:`heartbeat_if_in_activity`.
28+
29+
Kept so the existing call sites keep working; prefer the activity-named
30+
helper in new code.
31+
"""
32+
heartbeat_if_in_activity(heartbeat_name)
33+
34+
2035
def workflow_now_if_in_workflow() -> datetime | None:
2136
# Returns Temporal's deterministic workflow clock when called from inside a
2237
# workflow, otherwise None. Used to stamp messages with a monotonic

‎tests/lib/test_temporal_utils.py‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,18 @@
22

33
from __future__ import annotations
44

5+
from typing import Any
56
from datetime import datetime
67
from unittest.mock import patch
78

9+
from temporalio import activity
10+
from temporalio.testing import ActivityEnvironment
11+
812
from agentex.lib.utils import temporal as _temporal_mod
913
from agentex.lib.utils.temporal import (
1014
in_temporal_workflow,
15+
heartbeat_if_in_activity,
16+
heartbeat_if_in_workflow,
1117
workflow_now_if_in_workflow,
1218
)
1319

@@ -18,6 +24,45 @@ def test_in_temporal_workflow_returns_false_outside_workflow() -> None:
1824
assert in_temporal_workflow() is False
1925

2026

27+
async def test_heartbeat_if_in_activity_heartbeats_inside_activity() -> None:
28+
"""The helper must reach Temporal's heartbeat channel from inside an activity."""
29+
recorded: list[tuple[Any, ...]] = []
30+
31+
@activity.defn(name="heartbeating_activity")
32+
async def heartbeating_activity() -> None:
33+
heartbeat_if_in_activity("doing slow work")
34+
35+
env = ActivityEnvironment()
36+
env.on_heartbeat = lambda *details: recorded.append(details)
37+
38+
await env.run(heartbeating_activity)
39+
40+
assert recorded == [("doing slow work",)]
41+
42+
43+
async def test_heartbeat_if_in_workflow_alias_heartbeats_inside_activity() -> None:
44+
"""The legacy name stays wired to the same behaviour for existing call sites."""
45+
recorded: list[tuple[Any, ...]] = []
46+
47+
@activity.defn(name="legacy_heartbeating_activity")
48+
async def legacy_heartbeating_activity() -> None:
49+
heartbeat_if_in_workflow("doing slow work")
50+
51+
env = ActivityEnvironment()
52+
env.on_heartbeat = lambda *details: recorded.append(details)
53+
54+
await env.run(legacy_heartbeating_activity)
55+
56+
assert recorded == [("doing slow work",)]
57+
58+
59+
def test_heartbeat_if_in_activity_is_noop_outside_activity() -> None:
60+
"""Sync agents and workflow code call the same helpers and must not raise there."""
61+
assert activity.in_activity() is False
62+
heartbeat_if_in_activity("doing slow work")
63+
heartbeat_if_in_workflow("doing slow work")
64+
65+
2166
def test_workflow_now_if_in_workflow_returns_none_outside_workflow() -> None:
2267
assert workflow_now_if_in_workflow() is None
2368

0 commit comments

Comments
 (0)