Compare commits

...
1 Commits
Author SHA1 Message Date
bruno 3bb8e87ef2 fix: A42 terminé — client httpx partagé par boucle (v7.28.0)
FlowDeck CI / lint (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 0s
FlowDeck CI / docker (push) Canceled after 0s
- `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
2026-10-01 23:09:45 -04:00
23 changed files with 200 additions and 71 deletions
+29
View File
@@ -1,5 +1,34 @@
# Changelog - FlowDeck
## v7.28.0 (2026-10-01) — Audit : A42 TERMINÉ (client httpx partagé)
### Changed
- **`app/services/http_client.py`** : `async with shared_client(timeout=15)
as client:` remplace les **49 créations `async with httpx.AsyncClient(`**
réparties dans 14 fichiers (gitea ×21, providers ×11, calendar ×4,
automations ×3, …) — le pool de connexions est réutilisé au lieu d'être
recréé à chaque appel
- Cache **par (boucle d'event, kwargs)** en `WeakKeyDictionary` : un
`AsyncClient` n'est jamais partagé entre deux loops (le piège des tests :
« Event loop is closed ») — une boucle par test = un client propre,
collecté avec elle. Clé = kwargs triés, `repr()` pour les valeurs non
hashables (`headers=` dict)
- Context manager no-op à la sortie (pas de fermeture du client partagé) ;
`ponytail:` documenté : pas d'`aclose` explicite, plafond = pools non
fermés à la main (GC des sockets), upgrade = lifespan
- **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),
cloisonnement par kwargs, cloisonnement par boucle (2× `asyncio.run`)
- Stub webhooks : patch étendu à la fabrique `http_client.httpx` + purge du
cache (les tests patchaient `webhook_outbound.httpx`, contourné par la
fabrique partagée)
- Suite complète : **1091/1091**
## v7.27.0 (2026-10-01) — Audit : A20 phase 2 (CDN retiré, connect-src fermé)
### Changed
+2 -2
View File
File diff suppressed because one or more lines are too long
+1 -1
View File
@@ -1 +1 @@
7.27.0
7.28.0
+1 -1
View File
@@ -1,6 +1,6 @@
# WORKLOAD — FlowDeck Notion Clone
> **Début**: 2026-07-08 | **Version**: v7.27.0 (audit — A20 phase 2 : CDN retiré, connect-src fermé) | **Statut**: EN COURS 🔄
> **Début**: 2026-07-08 | **Version**: v7.28.0 (audit — A42 TERMINÉ : client httpx partagé par boucle) | **Statut**: EN COURS 🔄
> **Cible**: parité Notion + intégration forge · **Follow-ups v7.3 livrés**: sidebar teamspaces, notif `page.updated`, charts `number` + dashboards multi-DB, unfurl forge, UI Settings → Audit — voir `ROADMAP.md § v7.3.0`
## Avancement Global
+3 -4
View File
@@ -4,9 +4,8 @@ from __future__ import annotations
import logging
from urllib.parse import urlencode
import httpx
from app.config import settings
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -39,7 +38,7 @@ class GiteaOAuth:
async def exchange_code(self, code: str) -> dict | None:
"""Exchange authorization code for access token."""
try:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.post(
self.TOKEN_URL,
data={
@@ -61,7 +60,7 @@ class GiteaOAuth:
async def get_user(self, access_token: str) -> dict | None:
"""Get user info from Gitea API."""
try:
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
self.USER_URL,
headers={"Authorization": f"token {access_token}"},
+7 -7
View File
@@ -7,7 +7,7 @@ import time
from abc import ABC, abstractmethod
from urllib.parse import urlencode
import httpx
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -76,7 +76,7 @@ class GiteaProvider(OAuthProvider):
"grant_type": "authorization_code",
"redirect_uri": redirect_uri or self.redirect_uri,
}
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.post(url, json=data, headers={"Accept": "application/json"})
if r.status_code != 200:
logger.error("Gitea token exchange failed: %s", r.text)
@@ -85,7 +85,7 @@ class GiteaProvider(OAuthProvider):
async def get_user(self, access_token: str) -> dict | None:
url = f"{self.base}/api/v1/user"
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.get(url, headers={"Authorization": f"token {access_token}"})
if r.status_code != 200:
return None
@@ -100,7 +100,7 @@ class GiteaProvider(OAuthProvider):
async def list_repositories(self, access_token: str) -> list[dict]:
repos = []
async with httpx.AsyncClient(timeout=30) as client:
async with shared_client(timeout=30) as client:
for page in range(1, 6):
r = await client.get(
f"{self.base}/api/v1/user/repos",
@@ -158,7 +158,7 @@ class GitHubProvider(OAuthProvider):
)
async def exchange_code(self, code: str, redirect_uri: str | None = None) -> dict | None:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.post(
self.token_url,
data={
@@ -179,7 +179,7 @@ class GitHubProvider(OAuthProvider):
async def get_user(self, access_token: str) -> dict | None:
url = f"{self.api_url}/user"
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.get(
url,
headers={"Authorization": f"Bearer {access_token}", "Accept": "application/vnd.github.v3+json"},
@@ -197,7 +197,7 @@ class GitHubProvider(OAuthProvider):
async def list_repositories(self, access_token: str) -> list[dict]:
repos = []
async with httpx.AsyncClient(timeout=30) as client:
async with shared_client(timeout=30) as client:
for page in range(1, 6):
r = await client.get(
f"{self.api_url}/user/repos",
+4 -4
View File
@@ -15,7 +15,7 @@ import secrets
import time
import warnings
import httpx
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -59,7 +59,7 @@ async def discover(issuer_url: str) -> dict:
if hit and now - hit[0] < _DISCOVERY_TTL:
return hit[1]
try:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.get(url)
r.raise_for_status()
doc = r.json()
@@ -112,7 +112,7 @@ async def exchange_code(
if client_secret:
auth = (client_id, client_secret)
try:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.post(doc["token_endpoint"], data=data, auth=auth)
except Exception as err:
raise OIDCError(f"OIDC token request failed: {err}") from err
@@ -133,7 +133,7 @@ async def fetch_userinfo(doc: dict, access_token: str) -> dict:
if not endpoint or not access_token:
return {}
try:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.get(endpoint, headers={"Authorization": f"Bearer {access_token}"})
if r.status_code != 200:
return {}
+1 -1
View File
@@ -185,7 +185,7 @@ async def lifespan(_app: FastAPI):
app = FastAPI(
title="FlowDeck",
version="7.27.0",
version="7.28.0",
docs_url="/docs",
redoc_url="/redoc",
lifespan=lifespan,
+2 -2
View File
@@ -16,6 +16,7 @@ from app.routers.dashboard import _get_app_version
from app.routers.sidebar_config import get_sidebar_config_sync
from app.services.automations import fire_event, run_event_sync
from app.services.gitea_client import gitea
from app.services.http_client import shared_client
from app.services.permission_manager import PermissionManager
from app.services.publish import fire_published, fire_unpublished, publish, unpublish
@@ -1880,8 +1881,7 @@ async def _unfurl_repo(forge: str, owner: str, repo: str):
if token:
info = await GitHubAdapter(access_token=token).get_repo_info(owner, repo)
else:
import httpx
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
r = await client.get(
f"https://api.github.com/repos/{owner}/{repo}",
headers={"Accept": "application/vnd.github+json"},
+2 -2
View File
@@ -18,6 +18,7 @@ from fastapi.responses import JSONResponse
from app.auth.session import SessionManager
from app.db import get_conn
from app.services.api_v2_helpers import audit_log
from app.services.http_client import shared_client
router = APIRouter(tags=["scim"])
SCIM_SCHEMAS = ["urn:ietf:params:scim:schemas:core:2.0:User"]
@@ -246,7 +247,6 @@ def create_domain(request: Request, body: dict = Body(default={})):
@router.post("/api/v2/domain-claims/{domain_id}/verify")
async def verify_domain(domain_id: int, request: Request):
admin = _admin_session(request)
import httpx
with get_conn() as conn:
row = conn.execute("SELECT * FROM domain_claims WHERE id=?", (domain_id,)).fetchone()
if not row:
@@ -254,7 +254,7 @@ async def verify_domain(domain_id: int, request: Request):
claim = dict(row)
url = f"https://{claim['domain']}/.well-known/flowdeck-verify.txt"
try:
async with httpx.AsyncClient(timeout=10, follow_redirects=True) as client:
async with shared_client(timeout=10, follow_redirects=True) as client:
resp = await client.get(url)
ok = resp.status_code == 200 and claim["txt_token"] in (resp.text or "")
except Exception: # noqa: BLE001 — unreachable domain = not verified
+2 -2
View File
@@ -24,6 +24,7 @@ from app.auth.providers import oidc_provider, saml_provider
from app.auth.session import SessionManager
from app.services import sso_provisioning as sso
from app.services.api_v2_helpers import has_scope, resolve_bearer_token
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
router = APIRouter(tags=["sso"])
@@ -436,10 +437,9 @@ async def _fetch_jwks(doc: dict) -> dict:
url = doc.get("jwks_uri")
if not url:
raise oidc_provider.OIDCError("Discovery document has no jwks_uri")
import httpx
try:
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
r = await client.get(url)
r.raise_for_status()
data = r.json()
+4 -5
View File
@@ -26,10 +26,9 @@ import logging
import time
from datetime import UTC, datetime, timedelta
import httpx
from app.db import get_conn
from app.services import notifications
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -179,7 +178,7 @@ async def _run_action(action: dict, context: dict, trigger_source: str) -> str:
headers = {"Content-Type": "application/json", "X-FlowDeck-Event": context.get("event", "")}
if secret:
headers["X-FlowDeck-Secret"] = secret
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.post(url, json=context, headers=headers)
if resp.status_code >= 400:
raise RuntimeError(f"webhook returned HTTP {resp.status_code}")
@@ -719,7 +718,7 @@ async def fire_stepped_event(event: str, payload: dict) -> None:
# ── v7.0.0 action backends (module-level = monkeypatchable in tests) ───────
async def _post_slack(webhook_url: str, text: str) -> str:
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.post(webhook_url, json={"text": text})
if resp.status_code >= 400:
raise RuntimeError(f"slack webhook returned HTTP {resp.status_code}")
@@ -765,7 +764,7 @@ async def _create_forge_issue(provider: str, owner: str, repo: str, title: str,
if provider == "github":
if not token:
raise ValueError("github action needs a linked GitHub account (token)")
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.post(
f"https://api.github.com/repos/{owner}/{repo}/issues",
headers={"Authorization": f"Bearer {token}",
+5 -6
View File
@@ -19,9 +19,8 @@ import time
import uuid
from datetime import UTC, datetime, timedelta
import httpx
from app.db import get_conn
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -125,7 +124,7 @@ async def google_list_events(tokens: dict, calendar_id: str,
url = (f"https://www.googleapis.com/calendar/v3/calendars/{calendar_id}"
f"/events?singleEvents=true&orderBy=startTime"
f"&timeMin={time_min}&timeMax={time_max}")
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(url, headers={"Authorization": f"Bearer {access}"})
if resp.status_code == 401:
raise SyncError("google token expired — relink the calendar")
@@ -150,7 +149,7 @@ async def google_push_event(tokens: dict, calendar_id: str, event: dict,
"start": {"date": event.get("start", "")[:10]},
"end": {"date": event.get("start", "")[:10]}}
base = f"https://www.googleapis.com/calendar/v3/calendars/{calendar_id}/events"
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
if remote_id:
resp = await client.patch(f"{base}/{remote_id}",
headers={"Authorization": f"Bearer {access}"}, json=body)
@@ -224,7 +223,7 @@ async def caldav_list_events(creds: dict, time_min: str, time_max: str) -> list[
body = _CALDAV_REPORT.format(
start=time_min.replace("-", "").split("T")[0] + "T000000Z",
end=time_max.replace("-", "").split("T")[0] + "T000000Z")
async with httpx.AsyncClient(timeout=15, auth=auth if auth[0] else None) as client:
async with shared_client(timeout=15, auth=auth if auth[0] else None) as client:
resp = await client.request("REPORT", url, content=body,
headers={"Depth": "1",
"Content-Type": "application/xml"})
@@ -245,7 +244,7 @@ async def caldav_push_event(creds: dict, event: dict, remote_id: str = "") -> st
remote_id if remote_id.startswith("http") else f"{url}/{remote_id}")
ics = _event_to_ics(uid.split("@")[0], event.get("title", ""),
event.get("start", ""), event.get("description", ""))
async with httpx.AsyncClient(timeout=15, auth=auth if auth[0] else None) as client:
async with shared_client(timeout=15, auth=auth if auth[0] else None) as client:
resp = await client.put(href, content=ics, headers={"Content-Type": "text/calendar"})
if resp.status_code >= 400:
raise SyncError(f"caldav returned HTTP {resp.status_code}")
+22 -23
View File
@@ -5,9 +5,8 @@ import logging
from datetime import datetime, timedelta
from typing import Any
import httpx
from app.config import settings
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -48,7 +47,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/user/repos",
headers=self._headers,
@@ -65,7 +64,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/orgs/{org}/repos",
headers=self._headers,
@@ -82,7 +81,7 @@ class GiteaClient:
cached = self._cached(cache_key)
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}",
headers=self._headers,
@@ -109,7 +108,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/user/orgs", headers=self._headers
)
@@ -128,7 +127,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/issues",
headers=self._headers,
@@ -140,7 +139,7 @@ class GiteaClient:
return data
async def get_issue(self, owner: str, repo: str, issue_id: int) -> dict:
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}",
headers=self._headers,
@@ -163,7 +162,7 @@ class GiteaClient:
payload["assignees"] = [assignee]
self._invalidate_issue_cache(owner, repo)
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.post(
f"{self._base}/repos/{owner}/{repo}/issues",
headers=self._headers,
@@ -177,7 +176,7 @@ class GiteaClient:
) -> dict:
"""Update an issue (title, body, state, labels, milestone, assignees)."""
self._invalidate_issue_cache(owner, repo)
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.patch(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}",
headers=self._headers,
@@ -191,7 +190,7 @@ class GiteaClient:
) -> list[dict]:
"""Replace all labels on an issue (PUT endpoint)."""
self._invalidate_issue_cache(owner, repo)
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.put(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}/labels",
headers=self._headers,
@@ -210,7 +209,7 @@ class GiteaClient:
self, owner: str, repo: str, issue_id: int
) -> list[dict]:
"""Get comments for an issue."""
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}/comments",
headers=self._headers,
@@ -232,7 +231,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/labels",
headers=self._headers,
@@ -252,7 +251,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/milestones",
headers=self._headers,
@@ -267,7 +266,7 @@ class GiteaClient:
async def list_webhooks(self, owner: str, repo: str) -> list[dict]:
"""List webhooks for a repo."""
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/hooks",
headers=self._headers,
@@ -292,7 +291,7 @@ class GiteaClient:
"events": events,
"active": True,
}
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.post(
f"{self._base}/repos/{owner}/{repo}/hooks",
headers=self._headers,
@@ -303,7 +302,7 @@ class GiteaClient:
async def delete_webhook(self, owner: str, repo: str, hook_id: int) -> bool:
"""Delete a webhook."""
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.delete(
f"{self._base}/repos/{owner}/{repo}/hooks/{hook_id}",
headers=self._headers,
@@ -319,7 +318,7 @@ class GiteaClient:
if cached:
return cached
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/collaborators",
headers=self._headers,
@@ -341,7 +340,7 @@ class GiteaClient:
cached = self._cached(cache_key)
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(url, headers=self._headers)
resp.raise_for_status()
data = resp.json()
@@ -357,7 +356,7 @@ class GiteaClient:
cached = self._cached(cache_key)
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
@@ -382,7 +381,7 @@ class GiteaClient:
}
if sha:
body["sha"] = sha
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.put(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
@@ -398,7 +397,7 @@ class GiteaClient:
async def delete_file(self, owner, repo, path, sha, message):
"""Delete a file from the repo."""
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.delete(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
@@ -422,7 +421,7 @@ class GiteaClient:
cached = self._cached(key)
if cached:
return cached
async with httpx.AsyncClient(timeout=15) as client:
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/commits",
headers=self._headers,
+68
View File
@@ -0,0 +1,68 @@
"""Client HTTP partagé (A42 — reliquat du audit).
51 créations `httpx.AsyncClient(...)` éparpillées dans 15 fichiers = un
nouveau pool de connexions par appel (pas de keep-alive). Ici le pool est
réutilisé, **closé par boucle d'event** : un `AsyncClient` ne traverse pas
de boucle à boucle (les tests en créent une par test → client propre par
boucle, collecté avec elle).
`async with shared_client(timeout=15) as client:` — même syntaxe que
before, l'`__aexit__` est un no-op (on ne ferme pas le client partagé).
ponytail: pas d'`aclose` explicite — le client meurt avec sa boucle
(WeakKeyDictionary, les sockets sont fermés par le GC) ; plafond : le pool
du loop produit n'est pas fermé à la main. Upgrade si besoin : lifespan
qui ferme le client du loop principal.
"""
from __future__ import annotations
import asyncio
import weakref
import httpx
_clients: weakref.WeakKeyDictionary[asyncio.AbstractEventLoop, dict] = (
weakref.WeakKeyDictionary()
)
def _hashable(v) -> bool:
try:
hash(v)
return True
except TypeError:
return False
def _key(kwargs: dict) -> tuple:
"""Clé de cache = kwargs (timeout/auth/transport/headers…). Valeurs non
hashables (dict `headers=`…) → repr, même contrat que la clé."""
return tuple(
(k, v if _hashable(v) else repr(v))
for k, v in sorted(kwargs.items())
)
def get_shared_client(**kwargs) -> httpx.AsyncClient:
loop = asyncio.get_running_loop()
per_loop = _clients.setdefault(loop, {})
k = _key(kwargs)
client = per_loop.get(k)
if client is None:
client = httpx.AsyncClient(**kwargs)
per_loop[k] = client
return client
class shared_client:
"""Context manager async : yield du client partagé, no-op à la sortie."""
def __init__(self, **kwargs):
self._kwargs = kwargs
async def __aenter__(self) -> httpx.AsyncClient:
return get_shared_client(**self._kwargs)
async def __aexit__(self, *exc) -> bool:
return False
+2 -1
View File
@@ -12,6 +12,7 @@ from urllib.parse import urlparse
import httpx
from app.services.export import markdown_to_blocks
from app.services.http_client import shared_client
from app.services.importers.base import ImportPage, ImportResult
from app.services.importers.html_notes import _html_to_markdown
@@ -61,7 +62,7 @@ async def fetch_url_result(url: str, *, transport: httpx.BaseTransport | None =
safe_url = _validate_url(url)
result = ImportResult(source="url")
try:
async with httpx.AsyncClient(
async with shared_client(
timeout=15, follow_redirects=True, transport=transport,
headers={"User-Agent": "FlowDeck-Importer/1.0"},
) as client:
+2 -1
View File
@@ -27,6 +27,7 @@ from dataclasses import dataclass, field
import httpx
from app.config import settings
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -270,7 +271,7 @@ class LLMClient:
if self.api_key:
headers["Authorization"] = f"Bearer {self.api_key}"
try:
async with httpx.AsyncClient(timeout=settings.agent_run_timeout_seconds) as client:
async with shared_client(timeout=settings.agent_run_timeout_seconds) as client:
resp = await client.post(self._endpoint(), json=payload, headers=headers)
resp.raise_for_status()
data = resp.json()
+3 -4
View File
@@ -14,6 +14,7 @@ import json
import time
from app.config import settings
from app.services.http_client import shared_client
from app.services.llm_client import PROVIDER_LABELS, PROVIDER_MODELS, PROVIDERS
# Providers whose /v1/models lists far more entries than /v1/chat/completions
@@ -331,7 +332,6 @@ async def fetch_provider_models(provider: str, *, api_key: str = "",
Anthropic's native model listing uses `x-api-key` + `anthropic-version`.
Returns a de-duplicated list capped at 300 models.
"""
import httpx
provider = provider.lower()
base = (api_base or "").strip() or (PROVIDERS.get(provider) or (None, None))[0]
@@ -345,7 +345,7 @@ async def fetch_provider_models(provider: str, *, api_key: str = "",
elif api_key:
headers = {"Authorization": f"Bearer {api_key}"}
async with httpx.AsyncClient(timeout=timeout) as client:
async with shared_client(timeout=timeout) as client:
resp = await client.get(f"{base_url}/models", headers=headers)
if resp.status_code >= 400:
body = (resp.text or "").strip()
@@ -405,7 +405,6 @@ async def _validate_chat_models(base_url: str, api_key: str, candidates: list[st
"""
import asyncio
import httpx
sem = asyncio.Semaphore(concurrency)
start = time.monotonic()
@@ -427,7 +426,7 @@ async def _validate_chat_models(base_url: str, api_key: str, candidates: list[st
return None
try:
async with sem:
async with httpx.AsyncClient(timeout=timeout) as client:
async with shared_client(timeout=timeout) as client:
resp = await client.post(url, headers=headers, json=payload)
code = resp.status_code
if code < 300 or code == 429:
+3 -2
View File
@@ -12,6 +12,8 @@ import logging
import re
from urllib.parse import urljoin, urlparse
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
_META_TAG_RE = re.compile(r"<meta\b[^>]*?>", re.I)
@@ -143,7 +145,6 @@ async def fetch_og_metadata(url: str, timeout: float = 6.0, transport=None) -> d
src = "https://" + src
base = {"url": src, "title": "", "description": "", "image": "", "site_name": "", "favicon": ""}
try:
import httpx
headers = {
"User-Agent": "FlowDeck/5.5 bookmark-fetcher (+https://flowdeck.dracodev.net)",
@@ -152,7 +153,7 @@ async def fetch_og_metadata(url: str, timeout: float = 6.0, transport=None) -> d
kwargs = {"timeout": timeout}
if transport is not None:
kwargs["transport"] = transport
async with httpx.AsyncClient(**kwargs) as client:
async with shared_client(**kwargs) as client:
resp = await _get_checked(client, src, headers)
except ValueError:
# A12 : hôte privé/loopback ou trop de redirections → refus explicite.
+3 -2
View File
@@ -15,6 +15,7 @@ import time
import httpx
from app.db import get_conn
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
@@ -189,7 +190,7 @@ async def fire_event(event: str, payload: dict):
if not subs:
logger.debug("No subscribers for event: %s", event)
return
async with httpx.AsyncClient(timeout=30.0) as client:
async with shared_client(timeout=30.0) as client:
for sub in subs:
try:
await deliver_to_sub(sub["id"], sub["url"], event,
@@ -219,7 +220,7 @@ async def retry_due_deliveries(now: float | None = None) -> int:
).fetchall()
if not rows:
return 0
async with httpx.AsyncClient(timeout=10) as client:
async with shared_client(timeout=10) as client:
for r in rows:
try:
payload = json.loads(r["payload"] or "{}")
+1 -1
View File
@@ -2,7 +2,7 @@
"openapi": "3.1.0",
"info": {
"title": "FlowDeck",
"version": "7.27.0"
"version": "7.28.0"
},
"paths": {
"/auth/register": {
+27
View File
@@ -309,6 +309,33 @@ def test_csp_no_cdn_and_vendor(client):
assert len(r.content) > 500, (path, len(r.content))
def test_http_client_shared_and_loop_scoped():
"""A42 : le client HTTP partagé est réutilisé dans la même boucle,
cloisonné par kwargs, et JAMAIS partagé entre deux boucles (un
AsyncClient lié à une boucle morte lèverait « Event loop is closed »)."""
import asyncio
from app.services.http_client import shared_client
async def same_loop():
async with shared_client(timeout=15) as a:
async with shared_client(timeout=15) as b:
assert a is b, "même boucle + mêmes kwargs = même client"
async with shared_client(timeout=30) as c:
assert c is not a, "kwargs différents = client différent"
return a
first = asyncio.run(same_loop())
# nouvelle boucle (façon tests : une boucle par test) → nouveau client
async def other_loop():
async with shared_client(timeout=15) as d:
assert d is not first, "client jamais réutilisé sur une boucle morte"
return d
second = asyncio.run(other_loop())
assert second is not first
def test_no_duplicate_routes():
"""A24 : deux routes même méthode+chemin → l'une écrase silencieusement l'autre."""
from app.main import app
+6
View File
@@ -357,6 +357,12 @@ def _patch_async_client(monkeypatch, handler) -> None:
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