mirror of
https://github.com/hoshikawa2/agent_platform_oci.git
synced 2026-09-07 18:23:46 +00:00
Problema no sequence do pubsub. Precisa ser por transaction_id somente
This commit is contained in:
@@ -65,18 +65,18 @@ def _legacy_agent_name() -> str:
|
|||||||
|
|
||||||
|
|
||||||
def _mongo_collection() -> str:
|
def _mongo_collection() -> str:
|
||||||
"""Return the MongoDB collection used for observer sequence counters.
|
"""Return the shared MongoDB collection used by every event producer.
|
||||||
|
|
||||||
TIM legacy deployments used an agent-specific collection name, commonly
|
The collection must not vary by agent. A transaction can emit GRL, AGA,
|
||||||
``{agent_name}_event_counters``. Keep an explicit env override for BO
|
NOC and other events from different components, and all of them must
|
||||||
environments that already provisioned the collection, and fall back to the
|
increment the same counter document. Deployments may override the name,
|
||||||
legacy naming convention when no collection is configured.
|
but the configured value must be identical in every producer/pod.
|
||||||
"""
|
"""
|
||||||
return (
|
return (
|
||||||
os.getenv("PUBSUB_SEQUENCE_MONGODB_COLLECTION")
|
os.getenv("PUBSUB_SEQUENCE_MONGODB_COLLECTION")
|
||||||
or os.getenv("MONGODB_EVENT_COUNTERS_COLLECTION")
|
or os.getenv("MONGODB_EVENT_COUNTERS_COLLECTION")
|
||||||
or os.getenv("EVENT_COUNTERS_COLLECTION")
|
or os.getenv("EVENT_COUNTERS_COLLECTION")
|
||||||
or f"{_legacy_agent_name()}_event_counters"
|
or "observer_event_counters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -89,7 +89,10 @@ def _ttl_seconds() -> int:
|
|||||||
|
|
||||||
|
|
||||||
def _fallback_enabled() -> bool:
|
def _fallback_enabled() -> bool:
|
||||||
return _env_bool("PUBSUB_SEQUENCE_MEMORY_FALLBACK", True)
|
# An in-memory fallback creates duplicate sequences when multiple pods or
|
||||||
|
# event producers handle the same transaction. Keep it opt-in only for
|
||||||
|
# local/single-process development.
|
||||||
|
return _env_bool("PUBSUB_SEQUENCE_MEMORY_FALLBACK", False)
|
||||||
|
|
||||||
|
|
||||||
def _key_prefix() -> str:
|
def _key_prefix() -> str:
|
||||||
@@ -106,18 +109,20 @@ def build_sequence_key(
|
|||||||
session_id: str | None,
|
session_id: str | None,
|
||||||
transaction_id: str | None = None,
|
transaction_id: str | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Build the counter key, preferring transaction-only isolation.
|
"""Build one counter key for the whole transaction.
|
||||||
|
|
||||||
``session_id`` remains as a compatibility fallback for older producers that
|
``agent_id`` is intentionally ignored for transaction-scoped counters.
|
||||||
do not yet send ``transactionId``. New events must use one counter per
|
A single transaction may emit events from different agents/components
|
||||||
transaction so concurrent requests in the same session do not share a
|
(for example GRL and AGA), and those events must share one monotonic
|
||||||
sequence. The agent is deliberately omitted from transaction-scoped keys:
|
sequence. ``session_id`` is retained only as a compatibility fallback when
|
||||||
one transaction can emit GRL, AGA, IC and NOC events through different
|
no transaction identifier is present.
|
||||||
framework components/agents, and all of them must share the same counter.
|
|
||||||
"""
|
"""
|
||||||
if transaction_id:
|
if transaction_id:
|
||||||
transaction = _safe_part(transaction_id, "unknown_transaction")
|
transaction = _safe_part(transaction_id, "unknown_transaction")
|
||||||
return f"{_key_prefix()}:transaction:{transaction}"
|
return f"{_key_prefix()}:transaction:{transaction}"
|
||||||
|
|
||||||
|
# Legacy fallback. Including the agent here avoids changing old session-only
|
||||||
|
# behavior, but new integrations should always provide transactionId.
|
||||||
agent = _safe_part(agent_id or os.getenv("AGENT_NAME"), "agent")
|
agent = _safe_part(agent_id or os.getenv("AGENT_NAME"), "agent")
|
||||||
session = _safe_part(session_id, "unknown_session")
|
session = _safe_part(session_id, "unknown_session")
|
||||||
return f"{_key_prefix()}:{agent}:session:{session}"
|
return f"{_key_prefix()}:{agent}:session:{session}"
|
||||||
@@ -267,10 +272,10 @@ async def next_sequence(
|
|||||||
) -> int | None:
|
) -> int | None:
|
||||||
"""Return the next observer sequence isolated by transaction.
|
"""Return the next observer sequence isolated by transaction.
|
||||||
|
|
||||||
The preferred scope is ``transaction_id`` only. ``agent_id + session_id``
|
The preferred scope is only ``transaction_id``. Agent/event family must
|
||||||
is used only as a backward-compatible fallback when the producer does not
|
never participate in the key because one transaction can emit events from
|
||||||
send a transaction identifier. Redis and MongoDB increments remain atomic
|
several components. ``session_id`` is used only as a backward-compatible
|
||||||
across Kubernetes replicas.
|
fallback. Redis and MongoDB increments remain atomic across replicas.
|
||||||
"""
|
"""
|
||||||
if not sequence_enabled() or (not transaction_id and not session_id):
|
if not sequence_enabled() or (not transaction_id and not session_id):
|
||||||
return None
|
return None
|
||||||
|
|||||||
Reference in New Issue
Block a user