Skip to content

[Bug] ADK 2.8.0: repeated HITL deadlocks on replay within one invocation (no LLM) #7027

Description

@Breaknus

Describe the Bug

On google-adk==2.8.0, a resumable Workflow using the documented review/revise HITL pattern deadlocks when approving after two rejected reviews. The error is a replay sequence barrier timeout on HITL@1.

This reproduces with only public FunctionNodes and InMemorySessionService: no LLM, network request, database, private API override, or SDK patch. The same App and Runner are used throughout. Interrupt IDs are unique, the review counter lives in native session state, and Runner infers the same invocation from the function responses.

The review function follows Feeding the answer back into a loop in the 2.8.0 HITL reference, including Event(output=response, route=..., state=...).

Steps to Reproduce

  1. Use Python 3.12 with google-adk==2.8.0 installed. No API key or model configuration is needed.
  2. Save the standalone code below as repro.py and run python repro.py.
  3. The script answers the three successive native requests with approved=False, False, then True.
  4. The third resume fails after the native 15-second replay barrier timeout.

Expected Behavior

Three unique review requests are followed by exactly one execution of Publish. The final native review_count should be 3, all within one invocation.

Observed Behavior

Requests review_0, review_1, and review_2 are emitted. On approval, Publish is never executed, and review_count remains 2:

RuntimeError: Replay divergence detected: Timed out waiting for sequence key 'HITL@1' to be unblocked.

The relevant stack is _workflow.py:_run_loop -> _workflow.py:return_ctx -> _replay_sequence_barrier.py:wait. The retained node paths before failure, with the common DocumentedHitl@1/ prefix omitted, are:

HITL@1
HITL@1
Revise@1
HITL@2
HITL@1
HITL@2
Revise@2
HITL@3

All recorded events with an invocation ID belong to a single invocation. The script prints the IDs and paths so this can be checked independently.

Environment Details

  • ADK: 2.8.0
  • Python: 3.12.12
  • OS: macOS, arm64
  • google-genai: 2.22.0
  • pydantic: 2.13.4
  • Session service: InMemorySessionService
  • ResumabilityConfig(is_resumable=True)
  • HITL node: rerun_on_resume=True

Model Information

  • LiteLLM: No, not used by this reproduction.
  • Model: N/A; no LLM calls.

Minimal Reproduction Code

"""Check the documented HITL loop with unchanged public ADK."""

import asyncio
import json
from pydantic import BaseModel
from google.adk import Context, Event, Workflow
from google.adk.apps import App, ResumabilityConfig
from google.adk.events import RequestInput
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.adk.workflow import START, Edge, FunctionNode
from google.genai import types


class ApprovalSchema(BaseModel):
    """Describe the documented approval response."""

    approved: bool


async def main():
    """Run the documented loop using one persistent App and Runner."""
    published = []

    async def review(ctx: Context, node_input: object):
        """Run the review body from the pinned official HITL reference."""
        review_count = ctx.state.get("review_count", 0)
        interrupt_id = f"review_{review_count}"
        response = ctx.resume_inputs.get(interrupt_id)
        if response:
            yield Event(
                output=response,
                route="approved" if response.get("approved") else "rejected",
                state={"review_count": review_count + 1},
            )
            return
        yield RequestInput(
            interrupt_id=interrupt_id,
            message="Approve this plan?",
            response_schema=ApprovalSchema,
        )

    async def revise(ctx: Context, node_input: object):
        """Send the documented rejected route back to review."""
        return Event(route="review")

    async def publish(ctx: Context, node_input: object):
        """Record reaching the documented approved successor."""
        published.append(True)
        return "PUBLISHED"

    hitl = FunctionNode(name="HITL", func=review, rerun_on_resume=True)
    revision = FunctionNode(name="Revise", func=revise)
    publication = FunctionNode(name="Publish", func=publish)
    app = App(
        name="documented-hitl",
        root_agent=Workflow(
            name="DocumentedHitl",
            edges=[
                Edge(from_node=START, to_node=hitl),
                Edge(from_node=hitl, to_node=revision, route="rejected"),
                Edge(from_node=revision, to_node=hitl, route="review"),
                Edge(from_node=hitl, to_node=publication, route="approved"),
            ],
        ),
        resumability_config=ResumabilityConfig(is_resumable=True),
    )
    sessions = InMemorySessionService()
    await sessions.create_session(
        app_name=app.name, user_id="user", session_id="session"
    )
    runner = Runner(app=app, session_service=sessions)
    message = types.Content(role="user", parts=[types.Part(text="Review this plan")])
    request_ids = []
    try:
        for phase in range(4):
            events = [
                event
                async for event in runner.run_async(
                    user_id="user", session_id="session", new_message=message
                )
            ]
            calls = [
                part.function_call
                for event in events
                for part in (event.content.parts or [] if event.content else [])
                if part.function_call and part.function_call.name == "adk_request_input"
            ]
            print(
                json.dumps({"phase": phase, "requests": [c.id for c in calls]}),
                flush=True,
            )
            if phase < 3:
                assert len(calls) == 1
                call = calls[0]
                assert call.id not in request_ids
                request_ids.append(call.id)
                message = types.Content(
                    role="user",
                    parts=[
                        types.Part(
                            function_response=types.FunctionResponse(
                                id=call.id,
                                name="adk_request_input",
                                response={"approved": phase == 2},
                            )
                        )
                    ],
                )
        assert published == [True]
    except Exception as error:
        print(
            json.dumps(
                {
                    "exception_type": type(error).__name__,
                    "exception": str(error),
                    "published": len(published),
                }
            ),
            flush=True,
        )
        raise
    finally:
        session = await sessions.get_session(
            app_name=app.name, user_id="user", session_id="session"
        )
        print(
            json.dumps(
                {
                    "review_count": session.state.get("review_count"),
                    "invocations": sorted(
                        {e.invocation_id for e in session.events if e.invocation_id}
                    ),
                    "trace": [
                        e.node_info.path
                        for e in session.events
                        if e.node_info.path
                        and (
                            "/HITL@" in e.node_info.path
                            or "/Revise@" in e.node_info.path
                        )
                    ],
                }
            ),
            flush=True,
        )


asyncio.run(main())

Control experiment / limited workaround

Removing only output=response, from the review function's Event(...), while retaining route, the state update, is_resumable=True, and every other line, makes this exact script pass: exit 0, review_count=3, Publish exactly once. The original script exits 1 with the barrier timeout.

This is a narrow control and a possible public-API workaround when downstream code needs only routing/state. It is not a general workaround for nodes whose downstream consumers need their output, and it does not establish recovery of histories already containing the duplicate events.

Possible Cause

Source inspection suggests an interaction between two native behaviors:

The graph starts by fast-forwarding HITL@1, but its barrier is waiting for Revise@1, which cannot be scheduled until HITL completes. This matches the observed trace. I am reporting the behavior rather than proposing a general SDK patch: nested/dynamic ordering would need to be considered when choosing a fix.

Related issue / Regression

This appears distinct from #6497: that issue involved events from an earlier invocation. This reproduction uses one invocation, one Runner, unique interrupt IDs and no LLM. The current invocation filter is present in 2.8.0.

Earlier ADK versions have not been tested, so I am not claiming a version regression. The failure was reproduced in both the original small split-node graph and this documented loop; the exact standalone script above was rerun before filing.

Activity

  1. chelsealong commented on Sep 5, 2026

    @chelsealong
    Contributor

    Picking this one up now — opening a PR shortly. Flagging it here so nobody duplicates the work; if someone is already on it, say so and I will drop mine.

  2. added theissue type on Sep 7, 2026
  3. added
    agent engine[Component] This issue is related to Vertex AI Agent Engine
    on Sep 7, 2026
  4. markns commented on Sep 7, 2026

    @markns

    Confirming this on 2.8.0 and on main (b0180620f), and adding a second shape that #7028 as currently written does not cover.

    Second shape: LlmAgent nodes with output_schema in a clarify loop

    Our workflow walks a product down three levels. Each level is an LlmAgent node with an output_schema (a "resolved" pick or a "needs clarification" question), followed by a rerun_on_resume=True function node that either routes onward or yields RequestInput and, on the answer, routes back to the classifier:

    START -> classify_a -> resolve_a --reclassify--> classify_a
                                    --resolved----> classify_b -> resolve_b --reclassify--> classify_b
                                                                            --resolved----> classify_c -> resolve_c -> done
    

    The third resume deadlocks:

    ADK as Workflow root as NodeTool under an LlmAgent
    2.8.0 Timed out waiting for sequence key 'classify_a@1' never resumes (separate issue, fixed by 6d1451806)
    main b0180620f 'classify_a@1' 'classify_a@1', surfaced as the tool's error result
    main + #7028 'classify_a@2' 'classify_a@2'
    main + #7028 + the change below completes completes, parent gets the FunctionResponse

    Why #7028 misses it

    #7028 fixes the position of a run id at its first genuine terminal event (output, route or error_code) and ignores later echoes. That works for the issue's HITL node, which yields Event(output=..., route=...) itself. An LlmAgent node with output_schema never persists an event carrying output or route of its own, so the first terminal event _scan_sequence sees for classify_a@2 is the echo from _maybe_reemit_replayed_output on the next resume. That echo lands after resolve_a@2's genuine route event, so the barrier expects resolve_a@2 before classify_a@2, which can never happen.

    Instrumented barrier sequence on main + #7028, third resume:

    [SEQ] ['classify_a@1', 'resolve_a@1', 'resolve_a@2', 'classify_a@2', 'classify_b@1', ...]
    [BARRIER wait classify_a@2] next_expected=resolve_a@2
    RuntimeError: Replay divergence detected: Timed out waiting for sequence key 'classify_a@2' to be unblocked.
    

    A change that fixes both shapes

    Fix each segment's position at its first terminal event and never move it. A node cannot legitimately move behind a successor that depends on it: the interrupted-then-completed case keeps its original position, and the re-emitted echo becomes a no-op for ordering. On top of #7028's _scan_sequence:

          if is_terminal_event(event):
            # First-seen ordering: a child's position is fixed by its FIRST
            # terminal event. Later terminal events for the same run id are either
            # the re-emitted replay echo or an interrupted node completing after
            # its resume; neither can move the node behind a dependent successor.
            if segment not in sequence:
              sequence.append(segment)

    With that, tests/unittests/workflow/utils/test_replay_manager.py (including the new test from #7028), test_node_tool.py, test_workflow_hitl.py and utils/test_rehydration_utils.py pass (134 passed, 1 xfailed), and the reproduction below completes as root and as a NodeTool.

    Reproduction (no API key; scripted BaseLlm)

    NODETOOL=1 wraps the same workflow as a NodeTool under an LlmAgent root.

    """Three-level clarify loop with LlmAgent nodes (output_schema) deadlocks on the 3rd resume."""
    
    from __future__ import annotations
    
    import asyncio
    import hashlib
    import json
    import os
    from typing import Any, AsyncGenerator, Generator, Literal
    
    from google.adk import Agent, Context, Event, Workflow
    from google.adk.apps import App, ResumabilityConfig
    from google.adk.events import RequestInput
    from google.adk.models.base_llm import BaseLlm
    from google.adk.models.llm_response import LlmResponse
    from google.adk.runners import InMemoryRunner
    from google.adk.workflow import node
    from google.genai import types
    from pydantic import BaseModel, model_validator
    
    
    class ClassifyArgs(BaseModel):
        description: str
    
        @model_validator(mode="before")
        @classmethod
        def _accept_bare_string(cls, data):
            return {"description": data} if isinstance(data, str) else data
    
    
    class Pick(BaseModel):
        status: Literal["resolved", "needs_clarification"]
        value: str | None = None
        question: str | None = None
    
    
    class PickLlm(BaseLlm):
        """Clarification on the first call, resolved on the second."""
    
        model: str = "scripted"
        level: str = ""
        calls: int = 0
    
        async def generate_content_async(self, llm_request, stream: bool = False):
            self.calls += 1
            body = (
                {"status": "needs_clarification", "question": f"Clarify {self.level}?"}
                if self.calls == 1
                else {"status": "resolved", "value": f"{self.level}-ok"}
            )
            yield LlmResponse(content=types.Content(role="model", parts=[types.Part(text=json.dumps(body))]))
    
    
    def make_level(level: str):
        classify = Agent(
            name=f"classify_{level}",
            model=PickLlm(level=level),
            instruction="Classify {description}",
            output_schema=Pick,
            output_key=f"pick_{level}",
        )
    
        @node(name=f"resolve_{level}", rerun_on_resume=True)
        def resolve(ctx: Context) -> Generator[Any, None, None]:
            desc = ctx.state.get("description", "")
            pick = ctx.state.get(f"pick_{level}") or {}
            if hasattr(pick, "model_dump"):
                pick = pick.model_dump()
            iid = f"clarify_{level}_{hashlib.blake2b(desc.encode(), digest_size=4).hexdigest()}"
            answer = ctx.resume_inputs.get(iid)
            if answer is None:
                if pick.get("status") == "resolved":
                    yield Event(route="resolved")
                    return
                yield RequestInput(interrupt_id=iid, message=pick.get("question", "?"))
                return
            yield Event(state={"description": f"{desc}; {level}={answer}"}, route="reclassify")
    
        return classify, resolve
    
    
    def seed(node_input: ClassifyArgs | str | dict | None) -> Generator[Any, None, None]:
        d = node_input.description if isinstance(node_input, ClassifyArgs) else (node_input.get("description") if isinstance(node_input, dict) else (node_input or ""))
        yield Event(state={"description": str(d)})
    
    
    def done(ctx: Context) -> Generator[Any, None, None]:
        yield Event(output=f"classified: {ctx.state.get('description')}")
    
    
    classify_a, resolve_a = make_level("a")
    classify_b, resolve_b = make_level("b")
    classify_c, resolve_c = make_level("c")
    
    classify = Workflow(
        name="classify",
        description="Three-level clarify loop.",
        input_schema=ClassifyArgs,
        edges=[
            ("START", seed, classify_a, resolve_a),
            (resolve_a, {"resolved": classify_b, "reclassify": classify_a}),
            (classify_b, resolve_b),
            (resolve_b, {"resolved": classify_c, "reclassify": classify_b}),
            (classify_c, resolve_c),
            (resolve_c, {"resolved": done, "reclassify": classify_c}),
        ],
    )
    
    
    class RootLlm(BaseLlm):
        model: str = "scripted"
        calls: int = 0
    
        async def generate_content_async(self, llm_request: Any, stream: bool = False) -> AsyncGenerator[LlmResponse, None]:
            self.calls += 1
            part = (
                types.Part(function_call=types.FunctionCall(id="fc-1", name="classify", args={"description": "widget"}))
                if self.calls == 1
                else types.Part(text="done.")
            )
            yield LlmResponse(content=types.Content(role="model", parts=[part]))
    
    
    if os.environ.get("NODETOOL") == "1":
        root = Agent(name="root", model=RootLlm(), instruction="Call classify.", tools=[classify])
    else:
        root = classify
    app = App(name="repro", root_agent=root, resumability_config=ResumabilityConfig(is_resumable=True))
    
    
    def _interrupt_ids(events):
        return [p.function_call.id for e in events for p in ((e.content.parts if e.content else None) or []) if p.function_call and p.function_call.name == "adk_request_input"]
    
    
    def _answer(fc_id: str) -> types.Content:
        return types.Content(role="user", parts=[types.Part(function_response=types.FunctionResponse(id=fc_id, name="adk_request_input", response={"result": "x"}))])
    
    
    async def main() -> None:
        print(f"root = {type(root).__name__} {root.name}")
        runner = InMemoryRunner(app=app, app_name="repro")
        session = await runner.session_service.create_session(app_name="repro", user_id="u")
    
        async def turn(msg, label):
            print(f"--- {label} ---")
            events = []
            try:
                async for ev in runner.run_async(user_id="u", session_id=session.id, new_message=msg):
                    events.append(ev)
                    for p in ((ev.content.parts if ev.content else None) or []):
                        if p.function_call:
                            print(f"  FC {p.function_call.name} id={p.function_call.id}")
                        elif p.function_response:
                            print(f"  FR {p.function_response.name} -> {str(p.function_response.response)[:60]}")
                    if isinstance(ev.output, str):
                        print(f"  OUT {ev.output!r}")
            except BaseException as exc:  # noqa: BLE001
                print(f"  !! {type(exc).__name__}: {str(exc)[:120]}")
            return events
    
        evs = await turn(types.Content(role="user", parts=[types.Part(text="classify a widget")]), "turn 1")
        n = 1
        while (ids := _interrupt_ids(evs)) and n < 6:
            n += 1
            evs = await turn(_answer(ids[-1]), f"turn {n}: answer {ids[-1]}")
    
    
    if __name__ == "__main__":
        asyncio.run(main())

    Output on main (b0180620f), root:

    --- turn 1 ---
      FC adk_request_input id=clarify_a_97c9093e
    --- turn 2: answer clarify_a_97c9093e ---
      FC adk_request_input id=clarify_b_5a1ebffa
    --- turn 3: answer clarify_b_5a1ebffa ---
      FC adk_request_input id=clarify_c_6ec1c453
    --- turn 4: answer clarify_c_6ec1c453 ---
      !! RuntimeError: Replay divergence detected: Timed out waiting for sequence key 'classify_a@1' to be unblocked.
    

    With main + #7028 + the first-seen change, turn 4 prints OUT 'classified: classify a widget; a=x; b=x; c=x' instead, and with NODETOOL=1 the parent additionally receives FR classify -> {'result': ...}.

    Happy to open a PR with the change and a test for this shape if that helps, or to fold it into #7028.

  5. surajksharma07 commented on Sep 7, 2026

    @surajksharma07
    Collaborator

    HITL@1 gets bumped behind Revise@1 once the fast-forward echo fires and the barrier deadlocks on the third resume, your root cause call was spot on. #7028's first-genuine-terminal-event fix is the right direction and clears the original script but markns's follow-up is real too: main + #7028 still hangs on the LlmAgent/output_schema shape since that node never has a genuine terminal event to anchor on. Tried the stricter first-seen-only tweak markns proposed on top of #7028's _scan_sequence and both repros pass clean. Worth folding that in before merge rather than closing this on the narrower fix alone.

  6. added a commit that references this issue on Sep 10, 2026
    4b819ab
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

agent engine[Component] This issue is related to Vertex AI Agent Engine

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions