Skip to content

Pub/Sub trigger: a redelivered message re-runs the agent in a new session, repeating tool side effects that already happened #7322

Description

@FurkhanShaikh

🔴 Required Information

Describe the Bug:
Suppose a Pub/Sub trigger run fails after a tool has already caused an external
side effect, such as a payment, an email or a ticket.

  • The endpoint returns HTTP 500 and Pub/Sub redelivers the message.
  • ADK runs the agent again from the start, in a new session with no record of
    the first run, so the side effect happens again.

Nothing inside the agent can prevent this. Pub/Sub keeps the messageId the
same across redeliveries ("A redelivered message retains the same message ID
between redelivery attempts"). ADK logs it but never passes it to the agent or
its tools, so a tool has nothing to deduplicate on.

In src/google/adk/cli/trigger_routes.py
(v2.10.0; byte-identical in 2.9.2):

  • L419:
    session_id = str(uuid.uuid4()), once per HTTP delivery.
  • L514-516:
    the agent receives {"data": ..., "attributes": ...}. The messageId is
    not included.
  • L518-522:
    the messageId is only logged.
  • L530-541:
    HTTP 500 for a TransientError after retries, and for any other exception.
  • L486,
    the endpoint's own description: "errors trigger Pub/Sub retry".

The ambient agents docs describe the
design: "Each redelivery creates a new session. Trigger workloads are stateless
by design." They don't say what this means for tools with side effects.

The in-process 429 retry is fine: it reuses the session, so the model can see
the tool call that already completed. The problem is the hand-off to Pub/Sub.

Steps to Reproduce:

  1. pip install google-adk==2.10.0
  2. Save the script under "Minimal Reproduction Code" below as repro.py.
  3. Run python repro.py.

It needs no network access and no API key.

  • The agent is a stock LlmAgent on ADK's Gemini class, with
    base_url pointed at a local fake Gemini endpoint.
  • The fake endpoint answers the first request with a pay_invoice call.
    It answers the request after the tool with HTTP 503 UNAVAILABLE ("The model
    is overloaded") once, then behaves normally.
  • The loop at the bottom plays Pub/Sub. After a non-2xx response it posts
    the same push envelope again, with the same messageId.

Expected Behavior:

One payment for one message. Failing that, a way for the agent or its tools to
recognise a redelivery, so they can skip a side effect that already happened.

Observed Behavior:

google-adk 2.10.0
delivery 1 of messageId 1234567890: HTTP 500
delivery 2 of messageId 1234567890: HTTP 200
payments made for one message: 2

For delivery 1, ADK logs
Error processing Pub/Sub message: 503 UNAVAILABLE. {'error': {'code': 503, 'status': 'UNAVAILABLE', 'message': 'The model is overloaded. Please try again later.'}}.

Environment Details:

  • ADK Library Version (pip show google-adk): 2.10.0. Also 2.9.2, whose
    trigger_routes.py is byte-identical.
  • Desktop OS: Windows 11
  • Python Version (python -V): 3.13.9

Model Information:

  • Are you using LiteLLM: No
  • Which model is being used: gemini-2.5-flash through ADK's Gemini class,
    answered by a local fake endpoint. No real model was called.

🟡 Optional Information

Regression:

No. The same code is in the commit that added the trigger endpoints
(e2d970f).

Additional Context:

The same duplicate has other routes. They were measured with the same kind of
local harness: a fake Gemini endpoint, and Pub/Sub emulated from its documented
redelivery rule.

What happens after the tool's side effect Deliveries Side effects
the next model call returns 503 500, 200 2
the model returns 429 through all 3 in-process retries 500, 200 2
the tool's own HTTP call times out after the provider committed 500, 200 2
no error: the run outlasts the subscription's acknowledgement deadline resent while the first run is still running 2 (both runs complete)
the failure repeats on every delivery 500 ×5 5

Two rows deserve a note:

  • The acknowledgement-deadline row needs no failure at all. A push
    subscription's acknowledgement deadline defaults to 10 seconds, and "if you
    send a negative acknowledgment or the acknowledgment deadline expires,
    Pub/Sub resends the message". The docs give "10 minutes (ack deadline)" as
    the maximum processing time, but don't mention that a subscription created
    with defaults gets 10 seconds.
  • The last row has no cap. Without a dead-letter policy, a failure that
    repeats after the side effect repeats the side effect on every redelivery,
    until the message expires. The docs recommend a dead-letter queue but don't
    say why it matters here.

Two related observations:

  • Eventarc (Python): direct CloudEvents do put ce-id in the agent's
    input (L617-627); Pub/Sub-wrapped events don't.

  • google/adk-go appears to share the design. I read its source but did
    not run it:

    • RetriableRunner.RunAgent creates one session per delivery (its comment:
      "One session per delivery").
    • messageContentFromPubSub leaves out the message ID.
    • A failed run returns 500.

    I'm happy to open a companion issue there.

Possible changes, roughly by cost:

  1. Pass the delivery's identity into the run. For example, put messageId
    (and ce-id) in session state or in the attributes the agent receives, so
    a tool can derive an idempotency key for its provider. This is small, and
    lets applications fix the problem themselves.
  2. Derive session_id from the delivery (e.g. uuid5 of subscription and
    messageId) instead of uuid4().
    • A redelivery then lands in the session that holds the completed tool
      calls, as the in-process retry already does.
    • ADK could also acknowledge a redelivery whose session already holds a
      final response, which covers the acknowledgement-deadline case.
    • A redelivery that arrives while the first run is still going would need
      a guard.
  3. Document it. Say that a redelivery repeats side effects that already
    happened. Recommend an acknowledgement deadline above the expected run time
    (up to 600 s), a dead-letter policy, and idempotent tools. Show how a
    publisher can carry an idempotency key in the message attributes, which do
    reach the agent today.

I'm happy to be told this is intended. In that case, (3) alone would still help.

Found while auditing retry behaviour across agent frameworks.

Minimal Reproduction Code:

"""Pub/Sub trigger: one message, two payments.

Self-contained: `pip install google-adk` then `python issue_repro.py`. No network
access and no API key; nothing leaves 127.0.0.1.

- The agent is a stock LlmAgent on ADK's stock Gemini class. `base_url` points
  at a local fake Gemini endpoint.
- The fake endpoint answers the first request with a `pay_invoice` call. It
  answers the request after the payment with HTTP 503 UNAVAILABLE ("The model
  is overloaded"), once, and then behaves normally.
- The loop at the bottom plays Pub/Sub: after a non-2xx response it posts the
  same push envelope again, with the same messageId.
"""

import base64
import json
import os
import sys
import tempfile
import textwrap
import threading
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path

GEMINI = ("127.0.0.1", 8944)
PAYMENTS = []  # the payment provider's ledger: one entry per irreversible charge
_calls = {"n": 0}


class FakeGemini(BaseHTTPRequestHandler):
    def do_POST(self):  # noqa: N802
        body = json.loads(self.rfile.read(int(self.headers["Content-Length"])))
        _calls["n"] += 1
        after_payment = "functionResponse" in body["contents"][-1]["parts"][-1]
        if after_payment and _calls["n"] == 2:
            return self._send(503, {"error": {"code": 503, "status": "UNAVAILABLE",
                                              "message": "The model is overloaded. Please try again later."}})
        part = ({"text": "Paid."} if after_payment else
                {"functionCall": {"name": "pay_invoice",
                                  "args": {"invoice": "INV-1", "amount": "100.00"}}})
        self._send(200, {"candidates": [{"content": {"role": "model", "parts": [part]},
                                         "finishReason": "STOP"}]})

    def _send(self, status, obj):
        data = json.dumps(obj).encode()
        self.send_response(status)
        self.send_header("Content-Type", "application/json")
        self.send_header("Content-Length", str(len(data)))
        self.end_headers()
        self.wfile.write(data)

    def log_message(self, *args):
        pass


AGENT = f'''
import __main__
from google.adk.agents import LlmAgent
from google.adk.models import Gemini

def pay_invoice(invoice: str, amount: str) -> dict:
    """Pay an invoice. Irreversible."""
    __main__.PAYMENTS.append((invoice, amount))
    return {{"status": "PAID", "payment_number": len(__main__.PAYMENTS)}}

root_agent = LlmAgent(
    name="payer",
    model=Gemini(model="gemini-2.5-flash", base_url="http://{GEMINI[0]}:{GEMINI[1]}"),
    instruction="Pay the invoice in the message once.",
    tools=[pay_invoice],
)
'''


def main() -> int:
    os.environ["GOOGLE_API_KEY"] = "placeholder"  # never sent anywhere but the fake
    os.environ.pop("GOOGLE_GENAI_USE_VERTEXAI", None)
    server = ThreadingHTTPServer(GEMINI, FakeGemini)
    threading.Thread(target=server.serve_forever, daemon=True).start()

    agents = Path(tempfile.mkdtemp())
    (agents / "payer").mkdir()
    (agents / "payer" / "__init__.py").write_text("from . import agent\n")
    (agents / "payer" / "agent.py").write_text(textwrap.dedent(AGENT))

    import google.adk
    from fastapi.testclient import TestClient
    from google.adk.cli.fast_api import get_fast_api_app

    app = get_fast_api_app(agents_dir=str(agents), web=False, trigger_sources=["pubsub"])
    client = TestClient(app)
    envelope = {"message": {"data": base64.b64encode(b'{"invoice": "INV-1", "amount": "100.00"}').decode(),
                            "messageId": "1234567890"},
                "subscription": "projects/p/subscriptions/invoices"}

    print(f"google-adk {google.adk.__version__}")
    for delivery in range(1, 4):
        r = client.post("/apps/payer/trigger/pubsub", json=envelope)
        print(f"delivery {delivery} of messageId 1234567890: HTTP {r.status_code}")
        if r.status_code in (102, 200, 201, 202, 204):
            break  # acknowledged; otherwise Pub/Sub redelivers the same message
    print(f"payments made for one message: {len(PAYMENTS)}")
    server.shutdown()
    return 0 if len(PAYMENTS) > 1 else 1


if __name__ == "__main__":
    sys.exit(main())

How often has this issue occurred?:

Every time, under the conditions above. The reproduction is deterministic.

Activity

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

Metadata

Metadata

Labels

No labels
No labels

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions