From 26d33892f3aa095fbcc785dcceb8813d7b2a4829 Mon Sep 17 00:00:00 2001 From: "cristiano.hoshikawa" Date: Wed, 29 Jul 2026 08:49:32 -0300 Subject: [PATCH] TransactionID Pubsub problem --- .../analytics/tim_payload_mapper.py | 1 + .../agent_framework/analytics/tim_sequence.py | 73 ++++++++++++------- 2 files changed, 48 insertions(+), 26 deletions(-) diff --git a/libs/agent_framework/src/agent_framework/analytics/tim_payload_mapper.py b/libs/agent_framework/src/agent_framework/analytics/tim_payload_mapper.py index 6de8b5f..c6c900f 100644 --- a/libs/agent_framework/src/agent_framework/analytics/tim_payload_mapper.py +++ b/libs/agent_framework/src/agent_framework/analytics/tim_payload_mapper.py @@ -63,6 +63,7 @@ def map_analytics_event_to_tim_flat_payload( # Tracking "eventType": event.get("eventType") or event_type, "traceId": _first(data, "traceId", "trace_id"), + "transactionId": _first(data, "transactionId", "transaction_id", "transactionID"), "spanId": _first(data, "spanId", "span_id"), "parentSpanId": _first(data, "parentSpanId", "parent_span_id"), "eventName": _first(data, "eventName", "name"), diff --git a/libs/agent_framework/src/agent_framework/analytics/tim_sequence.py b/libs/agent_framework/src/agent_framework/analytics/tim_sequence.py index 6e5d5c7..1d64c04 100644 --- a/libs/agent_framework/src/agent_framework/analytics/tim_sequence.py +++ b/libs/agent_framework/src/agent_framework/analytics/tim_sequence.py @@ -101,10 +101,24 @@ def _safe_part(value: Any, fallback: str) -> str: return text.replace(" ", "_").replace("/", "_").replace("\\", "_") -def build_sequence_key(agent_id: str | None, session_id: str) -> str: +def build_sequence_key( + agent_id: str | None, + session_id: str | None, + transaction_id: str | None = None, +) -> str: + """Build the counter key, preferring transaction isolation. + + ``session_id`` remains as a compatibility fallback for older producers that + do not yet send ``transactionId``. New events must use one counter per + transaction so concurrent requests in the same session do not share a + sequence. + """ agent = _safe_part(agent_id or os.getenv("AGENT_NAME"), "agent") + if transaction_id: + transaction = _safe_part(transaction_id, "unknown_transaction") + return f"{_key_prefix()}:{agent}:transaction:{transaction}" session = _safe_part(session_id, "unknown_session") - return f"{_key_prefix()}:{agent}:{session}" + return f"{_key_prefix()}:{agent}:session:{session}" async def _next_sequence_redis(key: str, ttl_seconds: int) -> int | None: @@ -137,7 +151,8 @@ _mongo_index_lock = asyncio.Lock() def _next_sequence_mongodb_sync( key: str, agent_id: str | None, - session_id: str, + session_id: str | None, + transaction_id: str | None, ttl_seconds: int, ) -> int | None: uri = _mongo_uri() @@ -157,6 +172,8 @@ def _next_sequence_mongodb_sync( "$set": { "agentId": agent_id or os.getenv("AGENT_NAME") or "agent", "sessionId": session_id, + "transactionId": transaction_id, + "sequenceScope": "transaction" if transaction_id else "session", "updatedAt": now, }, "$setOnInsert": { @@ -214,7 +231,8 @@ async def _ensure_mongo_ttl_index_once(ttl_seconds: int) -> None: async def _next_sequence_mongodb( key: str, agent_id: str | None, - session_id: str, + session_id: str | None, + transaction_id: str | None, ttl_seconds: int, ) -> int | None: if not _mongo_uri(): @@ -226,6 +244,7 @@ async def _next_sequence_mongodb( key, agent_id, session_id, + transaction_id, ttl_seconds, ) except Exception: @@ -239,33 +258,26 @@ async def _next_sequence_memory(key: str) -> int: return _memory_counters[key] -async def next_sequence(agent_id: str | None, session_id: str | None) -> int | None: - """Return the next per-agent/per-session observer sequence. +async def next_sequence( + agent_id: str | None, + session_id: str | None, + transaction_id: str | None = None, +) -> int | None: + """Return the next observer sequence isolated by transaction. - Shared backends: - - Redis: atomic INCR, selected by PUBSUB_SEQUENCE_PROVIDER=redis. - - MongoDB: atomic find_one_and_update/$inc, selected by - PUBSUB_SEQUENCE_PROVIDER=mongodb. This mirrors the TIM legacy behavior. - - Provider selection: - - auto (default): Redis when configured; otherwise MongoDB when configured; - otherwise memory fallback when enabled. - - redis: Redis only, then memory fallback when enabled. - - mongodb/mongo: MongoDB only, then memory fallback when enabled. - - memory: in-process only. - - none: disabled. - - If session_id is absent or the shared backend fails and memory fallback is - disabled, None is returned so the payload remains valid without sequence. + The preferred scope is ``agent_id + transaction_id``. ``session_id`` is + used only as a backward-compatible fallback when the producer does not send + a transaction identifier. Redis and MongoDB increments remain atomic across + Kubernetes replicas. """ - if not sequence_enabled() or not session_id: + if not sequence_enabled() or (not transaction_id and not session_id): return None provider = _sequence_provider() if provider == "none": return None - key = build_sequence_key(agent_id, session_id) + key = build_sequence_key(agent_id, session_id, transaction_id) ttl_seconds = _ttl_seconds() value: int | None = None @@ -275,12 +287,16 @@ async def next_sequence(agent_id: str | None, session_id: str | None) -> int | N if provider == "redis": value = await _next_sequence_redis(key, ttl_seconds) elif provider == "mongodb": - value = await _next_sequence_mongodb(key, agent_id, session_id, ttl_seconds) + value = await _next_sequence_mongodb( + key, agent_id, session_id, transaction_id, ttl_seconds + ) else: # auto if _redis_url(): value = await _next_sequence_redis(key, ttl_seconds) if value is None and _mongo_uri(): - value = await _next_sequence_mongodb(key, agent_id, session_id, ttl_seconds) + value = await _next_sequence_mongodb( + key, agent_id, session_id, transaction_id, ttl_seconds + ) if value is not None: return value @@ -296,8 +312,13 @@ async def ensure_sequence(payload: dict[str, Any]) -> dict[str, Any]: if payload.get("sequence") is not None: return payload session_id = payload.get("sessionId") or payload.get("session_id") + transaction_id = ( + payload.get("transactionId") + or payload.get("transaction_id") + or payload.get("transactionID") + ) agent_id = payload.get("agentId") or payload.get("agent_id") or os.getenv("AGENT_NAME") - seq = await next_sequence(agent_id, session_id) + seq = await next_sequence(agent_id, session_id, transaction_id) if seq is not None: payload["sequence"] = seq return payload