fix: A42 terminé — client httpx partagé par boucle (v7.28.0)
- `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
This commit is contained in:
@@ -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 "{}")
|
||||
|
||||
Reference in New Issue
Block a user