Files
flowdeck/app/services/gitea_client.py
T
bruno 3bb8e87ef2
FlowDeck CI / lint (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 0s
FlowDeck CI / docker (push) Canceled after 0s
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
2026-10-01 23:09:45 -04:00

458 lines
16 KiB
Python

"""FlowDeck — Gitea API client with caching."""
from __future__ import annotations
import logging
from datetime import datetime, timedelta
from typing import Any
from app.config import settings
from app.services.http_client import shared_client
logger = logging.getLogger(__name__)
class GiteaClient:
"""Async Gitea API v1 client with simple TTL cache.
Can be instantiated with a per-user token (OAuth) or falls back
to the server-wide admin token.
"""
def __init__(self, user_token: str | None = None) -> None:
self._base = f"{settings.gitea_url}/api/v1"
token = user_token or settings.gitea_token
self._headers = {"Authorization": f"token {token}"}
self._cache: dict[str, tuple[datetime, Any]] = {}
self._ttl = timedelta(seconds=settings.gitea_cache_ttl)
# ── cache helper ──
def _cached(self, key: str) -> Any | None:
entry = self._cache.get(key)
if entry and entry[0] > datetime.now():
return entry[1]
return None
def _set_cache(self, key: str, value: Any) -> None:
now = datetime.now()
# A42 : évacue les entrées expirées (le dict ne pouvait que grandir)
self._cache = {k: v for k, v in self._cache.items() if v[0] > now}
self._cache[key] = (now + self._ttl, value)
# ── repos ──
async def get_user_repos(self, page: int = 1, limit: int = 30) -> list[dict]:
cache_key = f"repos:{page}:{limit}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/user/repos",
headers=self._headers,
params={"page": page, "limit": limit, "sort": "updated", "order": "desc"},
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
async def get_org_repos(self, org: str, page: int = 1, limit: int = 20) -> list[dict]:
cache_key = f"org_repos:{org}:{page}:{limit}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/orgs/{org}/repos",
headers=self._headers,
params={"page": page, "limit": limit},
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
async def get_repo_info(self, owner: str, repo: str) -> dict:
"""Repository metadata for an unfurl card (owner/name/branch/…)."""
cache_key = f"repo_info:{owner}:{repo}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}",
headers=self._headers,
)
resp.raise_for_status()
info = resp.json()
repo_info = {
"id": info.get("id"),
"name": info.get("name"),
"owner": (info.get("owner") or {}).get("login", owner),
"full_name": info.get("full_name") or f"{owner}/{repo}",
"clone_url": info.get("clone_url", ""),
"html_url": info.get("html_url", ""),
"default_branch": info.get("default_branch", "main"),
"description": info.get("description") or "",
"language": info.get("language") or "",
}
self._set_cache(cache_key, repo_info)
return repo_info
async def get_user_orgs(self) -> list[dict]:
cache_key = "user_orgs"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/user/orgs", headers=self._headers
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
# ── issues ──
async def get_issues(
self, owner: str, repo: str, state: str = "all", page: int = 1, limit: int = 50
) -> list[dict]:
cache_key = f"issues:{owner}:{repo}:{state}:{page}:{limit}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/issues",
headers=self._headers,
params={"state": state, "page": page, "limit": limit},
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
async def get_issue(self, owner: str, repo: str, issue_id: int) -> dict:
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}",
headers=self._headers,
)
resp.raise_for_status()
return resp.json()
async def create_issue(
self, owner: str, repo: str, title: str, body: str = "",
labels: list[int] | None = None, milestone: int | None = None,
assignee: str = "",
) -> dict:
"""Create a new issue in Gitea."""
payload: dict[str, Any] = {"title": title, "body": body}
if labels:
payload["labels"] = labels
if milestone:
payload["milestone"] = milestone
if assignee:
payload["assignees"] = [assignee]
self._invalidate_issue_cache(owner, repo)
async with shared_client(timeout=10) as client:
resp = await client.post(
f"{self._base}/repos/{owner}/{repo}/issues",
headers=self._headers,
json=payload,
)
resp.raise_for_status()
return resp.json()
async def update_issue(
self, owner: str, repo: str, issue_id: int, **kwargs
) -> dict:
"""Update an issue (title, body, state, labels, milestone, assignees)."""
self._invalidate_issue_cache(owner, repo)
async with shared_client(timeout=10) as client:
resp = await client.patch(
f"{self._base}/repos/{owner}/{repo}/issues/{issue_id}",
headers=self._headers,
json=kwargs,
)
resp.raise_for_status()
return resp.json()
async def update_issue_labels(
self, owner: str, repo: str, issue_id: int, labels: list[int]
) -> list[dict]:
"""Replace all labels on an issue (PUT endpoint)."""
self._invalidate_issue_cache(owner, repo)
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,
json={"labels": labels},
)
resp.raise_for_status()
return resp.json()
async def close_issue(self, owner: str, repo: str, issue_id: int) -> dict:
return await self.update_issue(owner, repo, issue_id, state="closed")
async def reopen_issue(self, owner: str, repo: str, issue_id: int) -> dict:
return await self.update_issue(owner, repo, issue_id, state="open")
async def get_issue_comments(
self, owner: str, repo: str, issue_id: int
) -> list[dict]:
"""Get comments for an issue."""
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,
)
resp.raise_for_status()
return resp.json()
def _invalidate_issue_cache(self, owner: str, repo: str) -> None:
prefix = f"issues:{owner}:{repo}:"
to_delete = [k for k in self._cache if k.startswith(prefix)]
for k in to_delete:
del self._cache[k]
# ── labels ──
async def get_labels(self, owner: str, repo: str) -> list[dict]:
cache_key = f"labels:{owner}:{repo}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/labels",
headers=self._headers,
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
# ── milestones ──
async def get_milestones(
self, owner: str, repo: str, state: str = "open"
) -> list[dict]:
cache_key = f"milestones:{owner}:{repo}:{state}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/milestones",
headers=self._headers,
params={"state": state},
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
# ── webhooks ──
async def list_webhooks(self, owner: str, repo: str) -> list[dict]:
"""List webhooks for a repo."""
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/hooks",
headers=self._headers,
)
resp.raise_for_status()
return resp.json()
async def create_webhook(
self, owner: str, repo: str, url: str, secret: str,
events: list[str] | None = None,
) -> dict:
"""Create a webhook for a repo."""
if events is None:
events = ["issues", "pull_request", "repository", "issue_comment"]
payload = {
"type": "gitea",
"config": {
"url": url,
"content_type": "json",
"secret": secret,
},
"events": events,
"active": True,
}
async with shared_client(timeout=10) as client:
resp = await client.post(
f"{self._base}/repos/{owner}/{repo}/hooks",
headers=self._headers,
json=payload,
)
resp.raise_for_status()
return resp.json()
async def delete_webhook(self, owner: str, repo: str, hook_id: int) -> bool:
"""Delete a webhook."""
async with shared_client(timeout=10) as client:
resp = await client.delete(
f"{self._base}/repos/{owner}/{repo}/hooks/{hook_id}",
headers=self._headers,
)
return resp.status_code == 204
# ── assignees/collaborators ──
async def get_collaborators(self, owner: str, repo: str) -> list[dict]:
"""List repo collaborators."""
cache_key = f"collaborators:{owner}:{repo}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=10) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/collaborators",
headers=self._headers,
)
resp.raise_for_status()
data = resp.json()
self._set_cache(cache_key, data)
return data
# ── repo contents ──
async def get_repo_contents(self, owner, repo, path=""):
"""Get contents of a repo directory (lazy: one level)."""
url = f"{self._base}/repos/{owner}/{repo}/contents"
if path:
url += f"/{path}"
cache_key = f"contents:{owner}:{repo}:{path}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(url, headers=self._headers)
resp.raise_for_status()
data = resp.json()
if isinstance(data, dict):
data = [data]
self._set_cache(cache_key, data)
return data
async def get_file_content(self, owner, repo, path):
"""Get raw file content (decoded from base64)."""
import base64
cache_key = f"file:{owner}:{repo}:{path}"
cached = self._cached(cache_key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
)
resp.raise_for_status()
item = resp.json()
content = ""
if item.get("encoding") == "base64" and item.get("content"):
try:
content = base64.b64decode(item["content"]).decode("utf-8")
except Exception:
content = "[binary file]"
self._set_cache(cache_key, content)
return content
async def create_or_update_file(self, owner, repo, path, content, message, sha=None):
"""Create or update a file. Requires `sha` for updates."""
import base64
body = {
"content": base64.b64encode(content.encode("utf-8")).decode("utf-8"),
"message": message,
}
if sha:
body["sha"] = sha
async with shared_client(timeout=15) as client:
resp = await client.put(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
json=body,
)
resp.raise_for_status()
data = resp.json()
parent = "/".join(path.split("/")[:-1])
for p in (parent, path.rsplit("/", 1)[0] if "/" in path else ""):
self._cache.pop(f"contents:{owner}:{repo}:{p}", None)
self._cache.pop(f"file:{owner}:{repo}:{path}", None)
return data
async def delete_file(self, owner, repo, path, sha, message):
"""Delete a file from the repo."""
async with shared_client(timeout=15) as client:
resp = await client.delete(
f"{self._base}/repos/{owner}/{repo}/contents/{path}",
headers=self._headers,
json={"message": message, "sha": sha},
)
resp.raise_for_status()
parent = "/".join(path.split("/")[:-1])
for p in (parent, path.rsplit("/", 1)[0] if "/" in path else ""):
self._cache.pop(f"contents:{owner}:{repo}:{p}", None)
self._cache.pop(f"file:{owner}:{repo}:{path}", None)
return {}
def _invalidate_tree_cache(self, owner, repo):
keys = [k for k in self._cache if k.startswith(f"contents:{owner}:{repo}:") or k.startswith(f"file:{owner}:{repo}:")]
for k in keys:
del self._cache[k]
async def get_file_commits(self, owner: str, repo: str, path: str) -> list[dict]:
"""Get recent commits affecting a specific file (cached 5 min)."""
key = f"commits:{owner}:{repo}:{path}"
cached = self._cached(key)
if cached:
return cached
async with shared_client(timeout=15) as client:
resp = await client.get(
f"{self._base}/repos/{owner}/{repo}/commits",
headers=self._headers,
params={"path": path, "limit": 20},
)
if resp.status_code == 200:
data = resp.json()
self._set_cache(key, data)
return data
return []
# ── singleton ──
gitea = GiteaClient()
# ── Helper: get per-user GiteaClient ──
def get_user_gitea_client(request):
"""Return a GiteaClient authenticated with the user's OAuth token, or None."""
from app.auth.session import SessionManager
user = SessionManager.decode_session(request.cookies.get("flowdeck_session", ""))
if not user:
return None
from app.db import get_conn
with get_conn() as conn:
row = conn.execute(
"SELECT access_token FROM user_oauth_tokens WHERE user_id=? AND provider='gitea' ORDER BY updated_at DESC LIMIT 1",
(user["id"],),
).fetchone()
if not row:
return None
return GiteaClient(user_token=row["access_token"])