- `app/services/http_client.py` : `async with shared_client(timeout=15) as client:` remplace les 49 créations `async with httpx.AsyncClient(` de 14 fichiers (gitea ×21, providers oidc/oauth ×11, calendar ×4, automations ×3…) — le pool de connexions est réutilisé au lieu d'être recréé à chaque appel. __aexit__ no-op (le client partagé ne se ferme pas à la sortie). - Cache par (boucle d'event, kwargs) en WeakKeyDictionary : un AsyncClient n'est JAMAIS partagé entre deux loops (piège des tests « Event loop is closed ») — une boucle par test = client propre collecté avec la boucle. Clé = kwargs triés, repr() pour les valeurs non hashables (`headers=` dict → TypeError rattrapé par la suite). - Laissés délibérément : github_adapter (transport MockTransport injecté), webhook_outbound (client « own_client » fermé par la fonction). - Tests : `test_http_client_shared_and_loop_scoped` (réutilisation mêmes kwargs / cloisonné kwargs / cloisonné loop) ; le stub des webhooks patche aussi la fabrique `http_client.httpx` + purge du cache (avant : webhook_outbound.httpx patché mais la fabrique partagée créait un vrai client → réseau réel dans les tests). suite **1091/1091** · ruff OK · docs à jour
527 lines
18 KiB
Python
527 lines
18 KiB
Python
"""FlowDeck — v6.4.0 Webhooks v2 tests (HMAC, retries, +20 events)."""
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
import tempfile
|
|
import time
|
|
from unittest.mock import AsyncMock, MagicMock
|
|
|
|
import httpx
|
|
import pytest
|
|
|
|
from app.db import get_conn, init_db
|
|
from app.services import webhook_outbound
|
|
from app.services.webhook_outbound import (
|
|
EVENTS,
|
|
MAX_ATTEMPTS,
|
|
RETRY_DELAYS,
|
|
_deliver_once,
|
|
_event_matches,
|
|
deliver_to_sub,
|
|
fire_event,
|
|
retry_due_deliveries,
|
|
sign_payload,
|
|
verify_signature,
|
|
)
|
|
|
|
# ── Fixtures ────────────────────────────────────────────────────────────────
|
|
|
|
@pytest.fixture
|
|
def db():
|
|
"""Fresh isolated SQLite database with the full FlowDeck schema.
|
|
|
|
Creates ``webhook_subscriptions`` + ``webhook_deliveries`` (versioned
|
|
migrations) so delivery journal rows can be asserted, then restores the
|
|
previous settings — safe under ``pytest -n auto``.
|
|
"""
|
|
tmp = tempfile.NamedTemporaryFile(suffix=".db", delete=False)
|
|
db_path = tmp.name
|
|
tmp.close()
|
|
|
|
prev_env = os.environ.get("DATABASE_URL")
|
|
os.environ["DATABASE_URL"] = f"sqlite:///{db_path}"
|
|
|
|
import app.config
|
|
settings_obj = app.config.settings
|
|
prev_url = settings_obj.database_url
|
|
settings_obj.database_url = f"sqlite:///{db_path}"
|
|
|
|
init_db()
|
|
|
|
yield
|
|
|
|
settings_obj.database_url = prev_url
|
|
if prev_env is None:
|
|
os.environ.pop("DATABASE_URL", None)
|
|
else:
|
|
os.environ["DATABASE_URL"] = prev_env
|
|
for suffix in ("", "-wal", "-shm"):
|
|
try:
|
|
os.unlink(db_path + suffix)
|
|
except (PermissionError, OSError):
|
|
pass
|
|
|
|
|
|
@pytest.fixture
|
|
def sample_payload():
|
|
return {"page_id": 123, "title": "Test Page", "user_id": 456}
|
|
|
|
|
|
@pytest.fixture
|
|
def sample_secret():
|
|
return "test_webhook_secret_12345"
|
|
|
|
|
|
# ── HMAC Signature Tests ────────────────────────────────────────────────────
|
|
|
|
|
|
def test_sign_payload_returns_sha256_hex():
|
|
"""Test that sign_payload returns a valid HMAC-SHA256 signature."""
|
|
secret = "my_secret"
|
|
body = b'{"event": "page.created", "data": {}}'
|
|
signature = sign_payload(secret, body)
|
|
|
|
assert signature.startswith("sha256=")
|
|
hex_part = signature[7:] # Remove 'sha256=' prefix
|
|
assert len(hex_part) == 64 # SHA-256 hex is 64 chars
|
|
|
|
|
|
def test_sign_payload_consistency():
|
|
"""Test that the same input produces the same signature."""
|
|
secret = "my_secret"
|
|
body = b'{"event": "page.created"}'
|
|
|
|
sig1 = sign_payload(secret, body)
|
|
sig2 = sign_payload(secret, body)
|
|
|
|
assert sig1 == sig2
|
|
|
|
|
|
def test_verify_signature_valid():
|
|
"""Test that verify_signature returns True for valid signatures."""
|
|
secret = "my_secret"
|
|
body = b'{"event": "page.created", "data": {}}'
|
|
signature = sign_payload(secret, body)
|
|
|
|
assert verify_signature(secret, body, signature) is True
|
|
|
|
|
|
def test_verify_signature_invalid_secret():
|
|
"""Test that verify_signature returns False for wrong secret."""
|
|
secret = "my_secret"
|
|
wrong_secret = "wrong_secret"
|
|
body = b'{"event": "page.created"}'
|
|
signature = sign_payload(secret, body)
|
|
|
|
assert verify_signature(wrong_secret, body, signature) is False
|
|
|
|
|
|
def test_verify_signature_invalid_signature():
|
|
"""Test that verify_signature returns False for tampered signature."""
|
|
secret = "my_secret"
|
|
body = b'{"event": "page.created"}'
|
|
signature = "sha256=invalid_hex_value"
|
|
|
|
assert verify_signature(secret, body, signature) is False
|
|
|
|
|
|
def test_verify_signature_empty_secret():
|
|
"""Test that verify_signature returns False for empty secret."""
|
|
body = b'{"event": "page.created"}'
|
|
signature = "sha256=somehash"
|
|
|
|
assert verify_signature("", body, signature) is False
|
|
|
|
|
|
def test_verify_signature_empty_signature():
|
|
"""Test that verify_signature returns False for empty signature."""
|
|
secret = "my_secret"
|
|
body = b'{"event": "page.created"}'
|
|
|
|
assert verify_signature(secret, body, "") is False
|
|
|
|
|
|
def test_verify_signature_none_values():
|
|
"""Test that verify_signature handles None values."""
|
|
assert verify_signature(None, b'{}', "sha256=abc") is False
|
|
assert verify_signature("secret", None, "sha256=abc") is False
|
|
assert verify_signature("secret", b'{}', None) is False
|
|
|
|
|
|
# ── Event Matching Tests ────────────────────────────────────────────────────
|
|
|
|
|
|
def test_event_matches_exact():
|
|
"""Test exact event matching."""
|
|
assert _event_matches("page.created", "page.created") is True
|
|
assert _event_matches("page.created", "page.updated") is False
|
|
|
|
|
|
def test_event_matches_wildcard_all():
|
|
"""Test wildcard '*' matches all events."""
|
|
assert _event_matches("*", "page.created") is True
|
|
assert _event_matches("*", "collection.deleted") is True
|
|
assert _event_matches("all", "agent.run.finished") is True
|
|
|
|
|
|
def test_event_matches_wildcard_prefix():
|
|
"""Test wildcard prefix matching (e.g., 'page.*')."""
|
|
assert _event_matches("page.*", "page.created") is True
|
|
assert _event_matches("page.*", "page.updated") is True
|
|
assert _event_matches("page.*", "page.deleted") is True
|
|
assert _event_matches("page.*", "collection.created") is False
|
|
|
|
|
|
def test_event_matches_wildcard_suffix():
|
|
"""Test that suffix wildcards are not supported (only prefix)."""
|
|
assert _event_matches("*.created", "page.created") is False
|
|
assert _event_matches("*.created", "collection.created") is False
|
|
|
|
|
|
# ── Event Catalogue Tests ────────────────────────────────────────────────────
|
|
|
|
|
|
def test_events_list_has_50_plus_events():
|
|
"""Test that EVENTS list has at least 50 events."""
|
|
assert len(EVENTS) >= 50, f"Expected at least 50 events, got {len(EVENTS)}"
|
|
|
|
|
|
def test_events_list_contains_new_events():
|
|
"""Test that new events are included in the catalogue."""
|
|
new_events = [
|
|
"page.renamed", "page.duplicated", "page.archived",
|
|
"comment.updated", "comment.deleted",
|
|
"mention.resolved",
|
|
"collection.renamed", "collection.duplicated",
|
|
"collection.page.moved",
|
|
"collection.view.updated", "collection.view.deleted",
|
|
"collection.property.created", "collection.property.updated", "collection.property.deleted",
|
|
"sprint.deleted", "sprint.started", "sprint.completed", "sprint.canceled",
|
|
"share.created", "share.updated", "share.revoked",
|
|
"automation.failed", "automation.retrying",
|
|
"agent.run.started", "agent.run.failed",
|
|
"workspace.created", "workspace.updated", "workspace.deleted",
|
|
"workspace.member.added", "workspace.member.removed", "workspace.member.role_changed",
|
|
"file.uploaded", "file.deleted",
|
|
"import.started", "import.completed", "import.failed",
|
|
]
|
|
for event in new_events:
|
|
assert event in EVENTS, f"Event '{event}' not found in EVENTS list"
|
|
|
|
|
|
def test_events_list_contains_legacy_events():
|
|
"""Test that legacy events are still present."""
|
|
legacy_events = [
|
|
"page.created", "page.updated", "page.deleted", "page.moved",
|
|
"comment.added", "comment.resolved",
|
|
"collection.created", "collection.updated", "collection.deleted",
|
|
"sprint.created", "sprint.updated",
|
|
"favorite.added", "favorite.removed",
|
|
"automation.fired", "agent.run.finished",
|
|
"ping",
|
|
]
|
|
for event in legacy_events:
|
|
assert event in EVENTS, f"Legacy event '{event}' missing from EVENTS list"
|
|
|
|
|
|
# ── Retry Schedule Tests ────────────────────────────────────────────────────
|
|
|
|
|
|
def test_retry_delays_are_2_10_60():
|
|
"""Test that retry delays are 2s, 10s, 60s."""
|
|
assert RETRY_DELAYS == (2, 10, 60)
|
|
|
|
|
|
def test_max_attempts_is_4():
|
|
"""Test that MAX_ATTEMPTS is 4 (1 initial + 3 retries)."""
|
|
assert MAX_ATTEMPTS == 4
|
|
|
|
|
|
# ── Delivery Tests (Mocked) ────────────────────────────────────────────────
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_once_success():
|
|
"""Test successful delivery in _deliver_once."""
|
|
mock_client = AsyncMock()
|
|
mock_response = MagicMock()
|
|
mock_response.status_code = 200
|
|
mock_client.post.return_value = mock_response
|
|
|
|
url = "https://example.com/webhook"
|
|
event = "page.created"
|
|
payload = {"page_id": 123}
|
|
secret = "test_secret"
|
|
|
|
code, error = await _deliver_once(mock_client, url, event, payload, secret)
|
|
|
|
assert code == 200
|
|
assert error == ""
|
|
mock_client.post.assert_called_once()
|
|
call_args = mock_client.post.call_args
|
|
assert call_args.kwargs["headers"]["X-FlowDeck-Event"] == event
|
|
assert call_args.kwargs["headers"]["X-FlowDeck-Signature"].startswith("sha256=")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_once_http_error():
|
|
"""Test _deliver_once handles HTTP errors."""
|
|
mock_client = AsyncMock()
|
|
mock_response = MagicMock()
|
|
mock_response.status_code = 500
|
|
mock_client.post.return_value = mock_response
|
|
|
|
url = "https://example.com/webhook"
|
|
event = "page.created"
|
|
payload = {}
|
|
secret = "test_secret"
|
|
|
|
code, error = await _deliver_once(mock_client, url, event, payload, secret)
|
|
|
|
assert code == 500
|
|
assert "HTTP 500" in error
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_once_timeout():
|
|
"""Test _deliver_once handles timeout errors."""
|
|
mock_client = AsyncMock()
|
|
mock_client.post.side_effect = httpx.TimeoutException("Timeout after 30s")
|
|
|
|
url = "https://example.com/webhook"
|
|
event = "page.created"
|
|
payload = {}
|
|
secret = "test_secret"
|
|
|
|
code, error = await _deliver_once(mock_client, url, event, payload, secret)
|
|
|
|
assert code is None
|
|
assert "Timeout" in error
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_once_request_error():
|
|
"""Test _deliver_once handles request errors."""
|
|
mock_client = AsyncMock()
|
|
mock_client.post.side_effect = httpx.RequestError("Connection refused")
|
|
|
|
url = "https://example.com/webhook"
|
|
event = "page.created"
|
|
payload = {}
|
|
secret = "test_secret"
|
|
|
|
code, error = await _deliver_once(mock_client, url, event, payload, secret)
|
|
|
|
assert code is None
|
|
assert "Request error" in error
|
|
|
|
|
|
# ── Integration Tests (isolated DB + mocked HTTP) ───────────────────────────
|
|
|
|
|
|
def _add_sub(url: str = "https://receiver.example/hook", event: str = "*",
|
|
secret: str = "test_secret") -> int:
|
|
"""Insert an active subscription and return its id."""
|
|
with get_conn() as conn:
|
|
cur = conn.execute(
|
|
"INSERT INTO webhook_subscriptions (url, event, secret, active) VALUES (?, ?, ?, 1)",
|
|
(url, event, secret),
|
|
)
|
|
conn.commit()
|
|
return cur.lastrowid
|
|
|
|
|
|
def _journal() -> list[dict]:
|
|
"""All delivery-journal rows, oldest first."""
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"SELECT * FROM webhook_deliveries ORDER BY id"
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
|
|
def _patch_async_client(monkeypatch, handler) -> None:
|
|
"""Route every ``httpx.AsyncClient`` built by webhook_outbound to a MockTransport."""
|
|
real_client = httpx.AsyncClient
|
|
transport = httpx.MockTransport(handler)
|
|
|
|
def factory(*args, **kwargs):
|
|
kwargs.pop("transport", None)
|
|
return real_client(transport=transport, **kwargs)
|
|
|
|
from types import SimpleNamespace
|
|
stub = SimpleNamespace(
|
|
AsyncClient=factory,
|
|
TimeoutException=httpx.TimeoutException,
|
|
RequestError=httpx.RequestError,
|
|
)
|
|
monkeypatch.setattr(webhook_outbound, "httpx", stub)
|
|
# A42 : la fabrique partagée crée les clients — elle doit voir le stub
|
|
# (boucle neuve par test → aucun cache à purger, on purge par sécurité).
|
|
from app.services import http_client as _hc
|
|
|
|
monkeypatch.setattr(_hc, "httpx", stub)
|
|
_hc._clients.clear()
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_to_sub_retries_on_failure(db, monkeypatch):
|
|
"""Two failures then success → 3 attempts, journal retrying→retrying→delivered."""
|
|
monkeypatch.setattr(webhook_outbound, "RETRY_DELAYS", (0, 0, 0)) # no real sleeps
|
|
sub_id = _add_sub(url="https://receiver.example/retry", event="page.created")
|
|
|
|
mock_client = AsyncMock()
|
|
mock_client.post.side_effect = [
|
|
MagicMock(status_code=500),
|
|
MagicMock(status_code=500),
|
|
MagicMock(status_code=200),
|
|
]
|
|
|
|
ok = await deliver_to_sub(
|
|
sub_id, "https://receiver.example/retry", "page.created",
|
|
{"page_id": 1}, "test_secret", _client=mock_client,
|
|
)
|
|
|
|
assert ok is True
|
|
assert mock_client.post.call_count == 3
|
|
assert [r["status"] for r in _journal()] == ["retrying", "retrying", "delivered"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_deliver_to_sub_fails_after_max_attempts(db, monkeypatch):
|
|
"""Persistent 500 → exactly MAX_ATTEMPTS attempts, final journal row 'failed'."""
|
|
monkeypatch.setattr(webhook_outbound, "RETRY_DELAYS", (0, 0, 0))
|
|
sub_id = _add_sub(url="https://receiver.example/down")
|
|
|
|
mock_client = AsyncMock()
|
|
mock_client.post.return_value = MagicMock(status_code=500)
|
|
|
|
ok = await deliver_to_sub(
|
|
sub_id, "https://receiver.example/down", "page.created",
|
|
{}, "test_secret", _client=mock_client,
|
|
)
|
|
|
|
assert ok is False
|
|
assert mock_client.post.call_count == MAX_ATTEMPTS
|
|
journal = _journal()
|
|
assert len(journal) == MAX_ATTEMPTS # 3 'retrying' + 1 'failed'
|
|
assert journal[-1]["status"] == "failed"
|
|
assert journal[-1]["http_code"] == 500
|
|
assert journal[-1]["attempt"] == MAX_ATTEMPTS - 1
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fire_event_ignores_unknown_events(db, monkeypatch):
|
|
"""fire_event returns before touching subscribers/HTTP for unknown events."""
|
|
delivered = AsyncMock()
|
|
monkeypatch.setattr(webhook_outbound, "deliver_to_sub", delivered)
|
|
_add_sub(url="https://receiver.example/any", event="*")
|
|
|
|
await fire_event("nope.not_an_event", {"x": 1})
|
|
|
|
delivered.assert_not_awaited()
|
|
assert _journal() == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fire_event_no_subscribers(db, monkeypatch):
|
|
"""Valid event but empty webhook_subscriptions → no delivery attempt."""
|
|
delivered = AsyncMock()
|
|
monkeypatch.setattr(webhook_outbound, "deliver_to_sub", delivered)
|
|
|
|
await fire_event("page.created", {"page_id": 1})
|
|
|
|
delivered.assert_not_awaited()
|
|
assert _journal() == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_fire_event_delivers_to_subscribers(db, monkeypatch):
|
|
"""Wildcard fan-out: only the subscription matching page.* receives the event."""
|
|
matching = _add_sub(url="https://receiver.example/pages", event="page.*")
|
|
_add_sub(url="https://receiver.example/collections", event="collection.created")
|
|
delivered = AsyncMock()
|
|
monkeypatch.setattr(webhook_outbound, "deliver_to_sub", delivered)
|
|
|
|
await fire_event("page.created", {"page_id": 42})
|
|
|
|
delivered.assert_awaited_once()
|
|
args = delivered.await_args.args
|
|
assert args[0] == matching
|
|
assert args[1] == "https://receiver.example/pages"
|
|
assert args[2] == "page.created"
|
|
assert args[3] == {"page_id": 42}
|
|
assert args[4] == "test_secret"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_full_webhook_flow(db, monkeypatch):
|
|
"""End-to-end: subscription → fire_event → signed HTTP POST → 'delivered' journal."""
|
|
secret = "flow_secret_abc"
|
|
sub_id = _add_sub(url="https://receiver.example/hook", event="page.*", secret=secret)
|
|
|
|
seen: dict = {}
|
|
|
|
def handler(request: httpx.Request) -> httpx.Response:
|
|
seen["url"] = str(request.url)
|
|
seen["headers"] = dict(request.headers)
|
|
seen["body"] = bytes(request.content)
|
|
return httpx.Response(200, json={"ok": True})
|
|
|
|
_patch_async_client(monkeypatch, handler)
|
|
|
|
await fire_event("page.created", {"page_id": 7, "title": "Hello"})
|
|
|
|
# HTTP assertions: right URL, event header, verifiable HMAC over the raw body
|
|
assert seen["url"] == "https://receiver.example/hook"
|
|
assert seen["headers"]["x-flowdeck-event"] == "page.created"
|
|
assert verify_signature(secret, seen["body"],
|
|
seen["headers"]["x-flowdeck-signature"]) is True
|
|
payload = json.loads(seen["body"])
|
|
assert payload["event"] == "page.created"
|
|
assert payload["page_id"] == 7
|
|
assert payload["title"] == "Hello"
|
|
|
|
# Journal assertions
|
|
journal = _journal()
|
|
assert len(journal) == 1
|
|
assert journal[0]["webhook_id"] == sub_id
|
|
assert journal[0]["status"] == "delivered"
|
|
assert journal[0]["http_code"] == 200
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_retry_due_deliveries(db, monkeypatch):
|
|
"""A 'retrying' row past next_retry_at is re-fired, delivered, then superseded."""
|
|
monkeypatch.setattr(webhook_outbound, "RETRY_DELAYS", (0, 0, 0))
|
|
sub_id = _add_sub(url="https://receiver.example/retry-due", event="page.*")
|
|
with get_conn() as conn:
|
|
conn.execute(
|
|
"""INSERT INTO webhook_deliveries
|
|
(webhook_id, event, status, http_code, error, duration_ms,
|
|
attempt, payload, next_retry_at)
|
|
VALUES (?, 'page.created', 'retrying', 500, 'HTTP 500', 12, 0, ?, ?)""",
|
|
(sub_id, json.dumps({"page_id": 9}), time.time() - 1),
|
|
)
|
|
conn.commit()
|
|
|
|
hits: list = []
|
|
|
|
def handler(request: httpx.Request) -> httpx.Response:
|
|
hits.append(request)
|
|
return httpx.Response(200)
|
|
|
|
_patch_async_client(monkeypatch, handler)
|
|
|
|
retried = await retry_due_deliveries(now=time.time())
|
|
|
|
assert retried == 1
|
|
assert len(hits) == 1
|
|
assert hits[0].headers["x-flowdeck-event"] == "page.created"
|
|
journal = _journal()
|
|
assert len(journal) == 2
|
|
assert journal[0]["status"] == "superseded" # original due row
|
|
assert journal[1]["status"] == "delivered" # new attempt
|
|
assert journal[1]["http_code"] == 200
|