Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3bb8e87ef2 |
@@ -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
File diff suppressed because one or more lines are too long
+1
-1
@@ -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
@@ -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 +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",
|
||||
|
||||
@@ -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
@@ -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,
|
||||
|
||||
@@ -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
@@ -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
@@ -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()
|
||||
|
||||
@@ -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}",
|
||||
|
||||
@@ -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}")
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
@@ -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:
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 "{}")
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
"openapi": "3.1.0",
|
||||
"info": {
|
||||
"title": "FlowDeck",
|
||||
"version": "7.27.0"
|
||||
"version": "7.28.0"
|
||||
},
|
||||
"paths": {
|
||||
"/auth/register": {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user