Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
import json
import logging
from contextlib import asynccontextmanager
from typing import TYPE_CHECKING, Any, AsyncGenerator
from typing import TYPE_CHECKING, AsyncGenerator

import httpx
from fastapi import APIRouter, FastAPI, HTTPException, Query, Request
Expand All @@ -41,10 +41,6 @@
logger = logging.getLogger(__name__)


class WebhookResponse(BaseModel):
result: Any = Field(..., description="The agent's response to the webhook")


class ErrorResponse(BaseModel):
error: str = Field(..., description="Error message")
detail: str | None = Field(None, description="Detailed error information")
Expand Down Expand Up @@ -218,10 +214,6 @@ def create_webhook_router(
if interface.prompt:
compiled_prompt = compile_template(interface.prompt)

# Get signature configuration
signature = interface.signature
output_is_string = signature.output.type == "string"

# Get subscription configuration
subscription = interface.subscription
secret = subscription.secret
Expand Down Expand Up @@ -260,13 +252,20 @@ async def websub_verification(
return PlainTextResponse(content=hub_challenge)
raise HTTPException(status_code=404, detail="Invalid mode")

async def _run_agent_in_background(user_prompt: str) -> None:
try:
response = await agent.arun(user_prompt)
logger.debug(f"Agent response: {response}")
except Exception:
logger.exception("Agent execution error")

# Webhook receiver endpoint
@router.post(
path,
status_code=202,
responses={
400: {"model": ErrorResponse},
401: {"model": ErrorResponse},
500: {"model": ErrorResponse},
},
)
async def receive_webhook(request: Request) -> JSONResponse:
Expand Down Expand Up @@ -309,33 +308,10 @@ async def receive_webhook(request: Request) -> JSONResponse:
# Default: stringify the payload
user_prompt = json.dumps(payload, indent=2)

try:
# Run the agent
response = await agent.arun(user_prompt)
logger.debug(f"Agent response: {response}")

# Format response based on output schema
if output_is_string:
if not isinstance(response, str):
response = json.dumps(response)
return JSONResponse(content={"result": response})
else:
if isinstance(response, dict):
return JSONResponse(content=response)
elif isinstance(response, str):
try:
return JSONResponse(content=json.loads(response))
except json.JSONDecodeError:
return JSONResponse(content={"result": response})
else:
return JSONResponse(content={"result": response})
task = asyncio.create_task(_run_agent_in_background(user_prompt))
task.add_done_callback(log_task_exception)

except Exception as e:
logger.exception("Agent execution error")
raise HTTPException(
status_code=500,
detail="Internal server error",
) from e
return JSONResponse(status_code=202, content={"status": "accepted"})

return router

Expand Down
21 changes: 10 additions & 11 deletions python-interpreter/packages/afm-core/tests/test_webhook.py
Original file line number Diff line number Diff line change
Expand Up @@ -219,7 +219,7 @@ def test_health_endpoint(self, mock_webhook_agent: MagicMock) -> None:
assert response.status_code == 200
assert response.json()["status"] == "ok"

def test_webhook_processes_payload(self, mock_webhook_agent: MagicMock) -> None:
def test_webhook_accepts_payload(self, mock_webhook_agent: MagicMock) -> None:
app = create_webhook_app(
mock_webhook_agent,
auto_subscribe=False,
Expand All @@ -233,11 +233,9 @@ def test_webhook_processes_payload(self, mock_webhook_agent: MagicMock) -> None:
headers={"User-Agent": "TestClient/1.0"},
)

assert response.status_code == 200
assert response.status_code == 202
data = response.json()
assert "result" in data
# The template should have substituted the values
assert "Processed:" in data["result"]
assert data["status"] == "accepted"

def test_webhook_with_signature_verification(
self, mock_webhook_agent: MagicMock
Expand All @@ -264,7 +262,7 @@ def test_webhook_with_signature_verification(
},
)

assert response.status_code == 200
assert response.status_code == 202

def test_webhook_rejects_invalid_signature(
self, mock_webhook_agent: MagicMock
Expand Down Expand Up @@ -300,9 +298,9 @@ def test_webhook_without_template_uses_raw_payload(
json={"type": "notification", "message": "Hello"},
)

assert response.status_code == 200
assert response.status_code == 202
data = response.json()
assert "Raw payload:" in data["result"]
assert data["status"] == "accepted"

def test_webhook_invalid_json_returns_400(
self, mock_webhook_agent: MagicMock
Expand All @@ -323,7 +321,7 @@ def test_webhook_invalid_json_returns_400(
assert response.status_code == 400
assert "Invalid JSON" in response.json()["detail"]

def test_webhook_agent_error_returns_500(
def test_webhook_agent_error_still_returns_202(
self, mock_webhook_agent: MagicMock
) -> None:

Expand All @@ -344,8 +342,9 @@ async def failing_arun(input_data: str, session_id: str = "default") -> str:
json={"event": "test"},
)

assert response.status_code == 500
assert "Internal server error" in response.json()["detail"]
# Fire-and-forget: always returns 202, agent errors are logged in background
assert response.status_code == 202
assert response.json()["status"] == "accepted"


class TestWebSubVerification:
Expand Down
Loading
Loading