Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions migrations/versions/20260929090000_a2a_message_idempotency.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
"""Add durable idempotency records for CRM A2A message/send requests."""

from typing import Sequence, Union

from alembic import op
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql

revision: str = "20260929090000"
down_revision: Union[str, None] = "26a14ac7025d"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.create_table(
"evo_ai_a2a_idempotency_records",
sa.Column("id", postgresql.UUID(as_uuid=True), nullable=False),
sa.Column("agent_id", postgresql.UUID(as_uuid=True), nullable=False),
sa.Column("context_id", sa.String(length=255), nullable=False),
sa.Column("idempotency_key", sa.String(length=128), nullable=False),
sa.Column("request_hash", sa.String(length=64), nullable=False),
sa.Column("state", sa.String(length=16), nullable=False),
sa.Column("response", sa.JSON(), nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.text("now()"), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.text("now()"), nullable=False),
sa.CheckConstraint("state IN ('processing', 'completed')", name="check_a2a_idempotency_state"),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint(
"agent_id",
"context_id",
"idempotency_key",
name="uq_a2a_idempotency_agent_context_key",
),
)
op.create_index(
"ix_a2a_idempotency_created_at",
"evo_ai_a2a_idempotency_records",
["created_at"],
)


def downgrade() -> None:
op.drop_index("ix_a2a_idempotency_created_at", table_name="evo_ai_a2a_idempotency_records")
op.drop_table("evo_ai_a2a_idempotency_records")
209 changes: 177 additions & 32 deletions src/api/a2a_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,12 @@
from src.services.tools_service import tools_service
from src.services.mcp_server_service import get_mcp_server
from src.middleware.permissions import RequirePermission
from src.services.a2a_idempotency import (
claim_or_replay,
complete as complete_idempotency,
release as release_idempotency,
request_fingerprint,
)

logger = setup_logger(__name__)

Expand Down Expand Up @@ -865,20 +871,17 @@ async def handle_message_send(
# Extract message from params
message = params.get("message")
if not message:
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_400_BAD_REQUEST),
message="Invalid params",
return JSONResponse(
status_code=status.HTTP_400_BAD_REQUEST,
details={
content={
"jsonrpc": "2.0",
"id": request_id,
"error": {
"code": -32602,
"message": "Invalid params",
"data": {"missing": "message"},
},
}
},
)

# Extract configuration from params (A2A spec: configuration is optional)
Expand Down Expand Up @@ -991,13 +994,131 @@ async def handle_message_send(
if not text and files:
text = "Analyze the provided files"

logger.info(f"📝 Extracted text: {text}")
logger.info(f"📝 Extracted text length: {len(text)}")
logger.info(f"📎 Extracted files: {len(files)}")

# Generate IDs
task_id = str(uuid.uuid4())
context_id = params.get("contextId", str(uuid.uuid4()))
if not isinstance(context_id, str) or len(context_id) > 255:
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_400_BAD_REQUEST),
message="Invalid contextId",
status_code=status.HTTP_400_BAD_REQUEST,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32602,
"message": "contextId must be a string of at most 255 characters",
}},
)
message_id = message.get("messageId")
idempotency_key = params.get("idempotencyKey") or message_id
# CRM messages have a stable persisted UUID. Generic A2A callers that omit
# one retain legacy behavior rather than being deduplicated on random IDs.
cacheable = bool(idempotency_key) and not push_notification_config
if idempotency_key and (not isinstance(idempotency_key, str) or len(idempotency_key) > 128):
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_400_BAD_REQUEST),
message="Invalid idempotency key",
status_code=status.HTTP_400_BAD_REQUEST,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32602,
"message": "idempotency key must be a string of at most 128 characters",
}},
)
fingerprint = None
owns_idempotency_record = False

if cacheable:
# Tool execution can depend on metadata and caller identity, not only the
# visible message/file parts. Include every JSON-RPC param in the hash so
# a reused key with altered tool context is rejected rather than replayed.
fingerprint = request_fingerprint(params)
try:
claim_state, cached_response = claim_or_replay(
db,
agent_id=agent_id,
context_id=str(context_id),
idempotency_key=str(idempotency_key),
fingerprint=fingerprint,
)
except Exception as error:
db.rollback()
logger.error(
"A2A idempotency storage unavailable; refusing unprotected execution "
f"(agent={agent_id}, context={context_id}, error={type(error).__name__})"
)
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_503_SERVICE_UNAVAILABLE),
message="Request could not be safely deduplicated; retry shortly",
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32000,
"message": "Idempotency service unavailable",
}},
)

if claim_state == "replay":
logger.info(
f"A2A duplicate replayed agent={agent_id} context={context_id} "
f"message_id={idempotency_key}"
)
return JSONResponse(
content={"jsonrpc": "2.0", "id": request_id, "result": cached_response}
)
if claim_state == "conflict":
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_409_CONFLICT),
message="Idempotency key was already used for a different message",
status_code=status.HTTP_409_CONFLICT,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32602,
"message": "Idempotency key payload mismatch",
}},
)
if claim_state == "in_progress":
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_409_CONFLICT),
message="The same message is currently being processed; retry shortly",
status_code=status.HTTP_409_CONFLICT,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32000,
"message": "Request already in progress",
}},
)
if claim_state == "expired":
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_409_CONFLICT),
message="Previous processing stopped unexpectedly; operator review is required",
status_code=status.HTTP_409_CONFLICT,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32000,
"message": "Processing lease expired",
}},
)
if claim_state != "owner":
logger.error(
f"Unexpected A2A idempotency state {claim_state!r}; refusing execution "
f"(agent={agent_id}, context={context_id})"
)
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_503_SERVICE_UNAVAILABLE),
message="Request could not be safely deduplicated; retry shortly",
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
details={"jsonrpc": "2.0", "id": request_id, "error": {
"code": -32000,
"message": "Unknown idempotency claim state",
}},
)
owns_idempotency_record = True

agent_execution_started = False
try:
# Extract conversation history for context
logger.info(
Expand All @@ -1020,18 +1141,18 @@ async def handle_message_send(

# Log detailed combined history for debugging
for i, msg in enumerate(combined_history):
logger.info(f" History[{i}]: {msg['role']} - {msg['content'][:50]}...")
logger.debug(f" History[{i}]: {msg['role']} - {len(msg['content'])} chars")

# Execute agent with files - the ADK runner will handle session history automatically
logger.info(
f"🤖 Executing agent {agent_id} with message: {text} and {len(files)} files"
f"🤖 Executing agent {agent_id} with message length {len(text)} and {len(files)} files"
)
logger.info(
f"📚 ADK will provide session context automatically ({len(combined_history)} previous messages available)"
)

metadata = extract_metadata_from_request(params)
logger.info(f"📋 Extracted metadata: {metadata}")
logger.info(f"📋 Extracted metadata fields: {len(metadata)}")

# Extract userId from params (contact_id sent from Rails)
# contextId is the conversation UUID, but userId should be the contact UUID
Expand All @@ -1049,6 +1170,7 @@ async def handle_message_send(

# CRM-236: bounded by whoever is waiting. When the caller hangs up the
# run is cancelled instead of burning more of the provider's quota.
agent_execution_started = True
result = await run_unless_client_disconnects(
request,
run_agent(
Expand All @@ -1068,7 +1190,7 @@ async def handle_message_send(
)

final_response = result.get("final_response", "No response")
logger.info(f"✅ Agent response: {final_response}")
logger.info(f"✅ Agent response generated ({len(final_response)} chars)")

# Log what we're about to send to create_task_response
logger.info(
Expand All @@ -1095,6 +1217,16 @@ async def handle_message_send(
current_user_message,
)

if owns_idempotency_record:
complete_idempotency(
db,
agent_id=agent_id,
context_id=str(context_id),
idempotency_key=str(idempotency_key),
fingerprint=fingerprint,
response=task_response,
)

logger.info(
f"📦 Task response created with {len(task_response.get('artifacts', []))} artifacts"
)
Expand All @@ -1113,6 +1245,21 @@ async def handle_message_send(
)

except ClientGoneAway as e:
if owns_idempotency_record and not agent_execution_started:
try:
release_idempotency(
db,
agent_id=agent_id,
context_id=str(context_id),
idempotency_key=str(idempotency_key),
fingerprint=fingerprint,
)
except Exception as release_error:
db.rollback()
logger.error(
"Could not release cancelled A2A idempotency record "
f"(agent={agent_id}, error={type(release_error).__name__})"
)
# Nobody is waiting for this response — it is logged for the operator
# and returned only because the handler must return something.
logger.warning(f"⛔ {e}")
Expand All @@ -1135,17 +1282,29 @@ async def handle_message_send(
)

except Exception as e:
if owns_idempotency_record and not agent_execution_started:
try:
release_idempotency(
db,
agent_id=agent_id,
context_id=str(context_id),
idempotency_key=str(idempotency_key),
fingerprint=fingerprint,
)
except Exception as release_error:
db.rollback()
logger.error(
"Could not release failed A2A idempotency record "
f"(agent={agent_id}, error={type(release_error).__name__})"
)
# CRM-236: a quota exhaustion and a bug of ours both answered "500 Agent
# execution failed". Report the provider's condition when we recognise it.
provider_failure = classify_provider_error(e)
if provider_failure is not None:
log_provider_failure(provider_failure, agent_id)
return error_response(
request=request,
code=map_status_to_error_code(provider_failure.http_status),
message=provider_failure.message,
return JSONResponse(
status_code=provider_failure.http_status,
details={
content={
"jsonrpc": "2.0",
"id": request_id,
"error": {
Expand All @@ -1163,22 +1322,8 @@ async def handle_message_send(
# and the one most likely to carry a credential (`?key=AIza…`).
safe = redact_secrets(str(e))
logger.error(f"❌ Agent execution error: {safe}")
return error_response(
request=request,
code=map_status_to_error_code(status.HTTP_500_INTERNAL_SERVER_ERROR),
message="Agent execution failed",
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
details={
"jsonrpc": "2.0",
"id": request_id,
"error": {
"code": -32603,
"message": "Agent execution failed",
"data": {"error": safe},
},
}
)
return JSONResponse(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
content={
"jsonrpc": "2.0",
"id": request_id,
Expand All @@ -1187,7 +1332,7 @@ async def handle_message_send(
"message": "Agent execution failed",
"data": {"error": safe},
},
}
},
)


Expand Down
Loading