import asyncio import os import re import shlex from contextlib import asynccontextmanager from typing import Any from urllib.parse import urljoin, urlsplit import httpx GITEA_URL = os.getenv("GITEA_URL", "http://127.0.0.1:3000").rstrip("/") GITEA_TOKEN = os.getenv("GITEA_TOKEN", "") REVIEW_DIFF_MAX_BYTES = 64 * 1024 REVIEW_DIFF_MAX_LINES = 400 AVAILABLE_ISSUE_PAGE_CONCURRENCY = 3 _client: "GiteaTransport | None" = None class GiteaOverloadedError(RuntimeError): """Raised when bounded transport admission expires before capacity is available.""" class GiteaTransport: """Application-lifetime HTTP transport with single-flight concurrent reads.""" def __init__( self, *, max_concurrency: int = 8, admission_timeout: float = 0.25, **kwargs, ) -> None: max_concurrency = max(2, max_concurrency) admission_timeout = max(0.001, admission_timeout) self._http = httpx.AsyncClient(base_url=GITEA_URL, timeout=10, **kwargs) self.max_concurrency = max_concurrency self.admission_timeout = admission_timeout self._request_slots = asyncio.Semaphore(max(2, max_concurrency)) self._read_slots = asyncio.Semaphore(max(1, max_concurrency - 1)) self._reads: dict[tuple, asyncio.Task[httpx.Response]] = {} @property def is_closed(self) -> bool: return self._http.is_closed def _read_key(self, url: str, kwargs: dict) -> tuple: params = tuple(httpx.QueryParams(kwargs.get("params", {})).multi_items()) headers = tuple(sorted(httpx.Headers(kwargs.get("headers", {})).multi_items())) return url, params, headers async def _acquire(self, semaphore: asyncio.Semaphore) -> None: try: await asyncio.wait_for(semaphore.acquire(), timeout=self.admission_timeout) except TimeoutError as exc: raise GiteaOverloadedError("Gitea transport is saturated") from exc async def _perform_get(self, url: str, kwargs: dict) -> httpx.Response: await self._acquire(self._read_slots) try: await self._acquire(self._request_slots) try: return await self._http.get(url, **kwargs) finally: self._request_slots.release() finally: self._read_slots.release() async def get(self, url: str, **kwargs) -> httpx.Response: key = self._read_key(url, kwargs) task = self._reads.get(key) if task is None: task = asyncio.create_task(self._perform_get(url, kwargs)) self._reads[key] = task task.add_done_callback( lambda completed, request_key=key: ( self._reads.pop(request_key, None) if self._reads.get(request_key) is completed else None ) ) return await asyncio.shield(task) async def _mutate(self, method: str, url: str, kwargs: dict) -> httpx.Response: await self._acquire(self._request_slots) try: return await self._http.request(method, url, **kwargs) finally: self._request_slots.release() async def post(self, url: str, **kwargs) -> httpx.Response: return await self._mutate("POST", url, kwargs) async def patch(self, url: str, **kwargs) -> httpx.Response: return await self._mutate("PATCH", url, kwargs) async def put(self, url: str, **kwargs) -> httpx.Response: return await self._mutate("PUT", url, kwargs) async def delete(self, url: str, **kwargs) -> httpx.Response: return await self._mutate("DELETE", url, kwargs) @asynccontextmanager async def stream(self, method: str, url: str, **kwargs): await self._acquire(self._read_slots) try: await self._acquire(self._request_slots) try: async with self._http.stream(method, url, **kwargs) as response: yield response finally: self._request_slots.release() finally: self._read_slots.release() def __getattr__(self, name: str): return getattr(self._http, name) async def aclose(self) -> None: tasks = list(self._reads.values()) for task in tasks: if not task.done(): task.cancel() if tasks: await asyncio.gather(*tasks, return_exceptions=True) self._reads.clear() await self._http.aclose() class WorkItems(list[dict]): """A list-compatible first page carrying truthful per-stream totals.""" def __init__(self, items: list[dict], pagination: dict[str, dict]): super().__init__(items) self.pagination = pagination class RepositoryItems(list[dict]): """A list-compatible repository page carrying truthful pagination metadata.""" def __init__(self, items: list[dict], pagination: dict): super().__init__(items) self.pagination = pagination class WorkRouteUnavailableError(ValueError): """Raised when a shared route no longer belongs in the current user's queue.""" class StaleReviewError(ValueError): """Raised before mutation when a pull request head changed during review.""" class InvalidReviewCommentError(ValueError): """Raised before mutation when an inline comment cannot target this change.""" class StalePullError(ValueError): """Raised before merge when an assigned pull request head changed.""" class PullNotMergeableError(ValueError): """Raised before merge when current pull state or checks prohibit it.""" class IssueNotAvailableError(ValueError): """Raised before assignment when an issue is no longer open and unassigned.""" class IssueEditConflictError(ValueError): """Raised when an issue changed after the editor loaded it.""" class IssueDependencyInvalidError(ValueError): """Raised when a requested blocker relationship is not valid.""" def _auth() -> dict[str, str]: headers: dict[str, str] = {"Accept": "application/json"} if GITEA_TOKEN: headers["Authorization"] = f"token {GITEA_TOKEN}" return headers def start_client(**kwargs) -> GiteaTransport: """Create the application-lifetime Gitea transport.""" global _client kwargs.setdefault("max_concurrency", int(os.getenv("GITEA_MAX_CONCURRENCY", "8"))) kwargs.setdefault( "admission_timeout", float(os.getenv("GITEA_ADMISSION_TIMEOUT_SECONDS", "0.25")), ) _client = GiteaTransport(**kwargs) return _client def _get_client() -> GiteaTransport: if _client is None or _client.is_closed: return start_client() return _client async def stop_client() -> None: global _client if _client is not None and not _client.is_closed: await _client.aclose() async def fetch(path: str) -> Any: r = await _get_client().get(f"/api/v1/{path}", headers=_auth()) r.raise_for_status() return r.json() async def fetch_text(path: str, max_bytes: int) -> tuple[str, bool]: chunks: list[bytes] = [] size = 0 truncated = False async with _get_client().stream( "GET", f"/api/v1/{path}", headers={**_auth(), "Accept": "text/plain"} ) as response: response.raise_for_status() async for chunk in response.aiter_bytes(): remaining = max_bytes - size if len(chunk) > remaining: chunks.append(chunk[:remaining]) truncated = True break chunks.append(chunk) size += len(chunk) return b"".join(chunks).decode("utf-8", errors="replace"), truncated def _diff_previews(diff: str, stream_truncated: bool) -> dict[str, dict]: previews: dict[str, dict] = {} current: dict | None = None in_hunk = False remaining = REVIEW_DIFF_MAX_LINES for line in diff.splitlines(): if line.startswith("diff --git "): try: target = shlex.split(line)[3] filename = target[2:] if target.startswith("b/") else target except (IndexError, ValueError): current = None continue current = { "diff_lines": [], "diff_available": False, "diff_binary": False, "diff_truncated": stream_truncated, } previews[filename] = current in_hunk = False continue if current is None: continue if line.startswith("Binary files ") or line == "GIT binary patch": current["diff_binary"] = True in_hunk = False continue if line.startswith("@@"): in_hunk = True if in_hunk and not line.startswith("\\ No newline at end of file"): if remaining: current["diff_lines"].append(line) current["diff_available"] = True remaining -= 1 else: current["diff_truncated"] = True return previews async def current_user() -> dict: return await fetch("user") async def repo_page(page: int = 1, limit: int = 50) -> dict: """Load one bounded page of repositories available to the current user.""" response = await _get_client().get( "/api/v1/user/repos", headers=_auth(), params={"page": page, "limit": limit}, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea repository response was not a list") items = [item for item in payload if isinstance(item, dict)] try: total = max(len(items), int(response.headers.get("X-Total-Count", len(items)))) except (TypeError, ValueError): total = len(items) return { "items": items, "page": page, "total": total, "has_more": page * limit < total, } async def repos() -> RepositoryItems: result = await repo_page() return RepositoryItems( result["items"], {key: result[key] for key in ("page", "total", "has_more")}, ) async def search_repositories(query: str, limit: int = 20) -> list[dict]: """Search repositories visible to the authenticated Gitea user.""" response = await _get_client().get( "/api/v1/repos/search", headers=_auth(), params={"q": query.strip(), "limit": limit, "page": 1}, ) response.raise_for_status() payload = response.json() items = payload.get("data") if isinstance(payload, dict) else None if not isinstance(items, list): raise ValueError("Gitea repository search response was not a list") return [item for item in items if isinstance(item, dict)] async def repository_access(repository: str) -> dict | None: """Return a repository only when the authenticated user can access it.""" response = await _get_client().get( f"/api/v1/repos/{repository}", headers=_auth() ) if response.status_code == 404: return None response.raise_for_status() payload = response.json() return payload if isinstance(payload, dict) else None WORK_SEARCHES = { "issue": ("assigned=true", "issues", None), "pull": ("assigned=true", "pulls", "assigned_to_me"), "review": ("review_requested=true", "pulls", "review_requested"), } async def work_page(stream: str, page: int = 1, limit: int = 50) -> dict: """Load exactly one bounded My Work stream page with upstream totals.""" selector, item_type, reason = WORK_SEARCHES[stream] response = await _get_client().get( "/api/v1/repos/issues/search", headers=_auth(), params={ "state": "open", selector.split("=", 1)[0]: "true", "type": item_type, "limit": limit, "page": page, }, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea work search response was not a list") items = [item for item in payload if isinstance(item, dict)] if reason: items = [{**item, "work_reasons": [reason]} for item in items] try: total = max(len(items), int(response.headers.get("X-Total-Count", len(items)))) except (TypeError, ValueError): total = len(items) return { "stream": stream, "items": items, "page": page, "total": total, "has_more": page * limit < total, } def _normalize_global_search_item(item: Any, kind: str) -> dict | None: if not isinstance(item, dict): return None repository = item.get("repository") repository = repository if isinstance(repository, dict) else {} url = _safe_web_url(item.get("html_url")) if not ( isinstance(item.get("number"), int) and isinstance(item.get("title"), str) and isinstance(item.get("state"), str) and isinstance(repository.get("full_name"), str) and repository.get("full_name") and url ): return None return { "kind": kind, "repository": repository["full_name"], "number": item["number"], "title": item["title"], "state": item["state"], "url": url, } async def global_search(query: str, limit: int = 10) -> dict: """Search accessible issues and pulls concurrently with a bounded result set.""" async def load(item_type: str) -> Any: response = await _get_client().get( "/api/v1/repos/issues/search", headers=_auth(), params={"q": query, "type": item_type, "limit": limit, "page": 1}, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea global search response was not a list") return payload outcomes = await asyncio.gather( load("issues"), load("pulls"), return_exceptions=True ) for outcome in outcomes: if isinstance(outcome, asyncio.CancelledError): raise outcome if all(isinstance(outcome, BaseException) for outcome in outcomes): raise outcomes[0] results = [] seen: set[tuple[str, str, int]] = set() for outcome, kind in zip(outcomes, ("issue", "pull"), strict=True): if isinstance(outcome, BaseException): continue for item in outcome: normalized = _normalize_global_search_item(item, kind) if normalized is not None: identity = (kind, normalized["repository"], normalized["number"]) if identity in seen: continue seen.add(identity) results.append(normalized) return { "items": results[:limit], "partial": any(isinstance(outcome, BaseException) for outcome in outcomes), } async def work_preview(repository: str, kind: str, number: int) -> dict: """Load bounded, read-only context for a global search result.""" issue, user = await asyncio.gather( fetch(f"repos/{repository}/issues/{number}"), current_user() ) if not isinstance(issue, dict): raise ValueError("Gitea work preview response was not an object") labels_value = issue.get("labels") assignees_value = issue.get("assignees") labels = labels_value if isinstance(labels_value, list) else [] assignees = assignees_value if isinstance(assignees_value, list) else [] assignee_names = [ assignee["login"] for assignee in assignees if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) ] author = issue.get("user") author = author if isinstance(author, dict) else {} login = user.get("login") if isinstance(user, dict) else "" state = issue.get("state") if isinstance(issue.get("state"), str) else "" actual_kind = "pull" if isinstance(issue.get("pull_request"), dict) else "issue" return { "kind": actual_kind, "repository": repository, "number": number, "title": issue.get("title", "") if isinstance(issue.get("title"), str) else "", "body": issue.get("body", "") if isinstance(issue.get("body"), str) else "", "state": state, "author": author.get("login", "") if isinstance(author.get("login"), str) else "", "labels": [ label["name"] for label in labels if isinstance(label, dict) and isinstance(label.get("name"), str) ], "assignees": assignee_names, "url": _safe_web_url(issue.get("html_url")), "claimable": actual_kind == "issue" and state == "open" and not assignee_names, "reopenable": actual_kind == "issue" and state == "closed", "assigned_to_me": bool(login and login in assignee_names), } def _normalize_available_issue(item: Any) -> dict | None: if ( not isinstance(item, dict) or item.get("state") != "open" or item.get("pull_request") is not None or item.get("assignees") not in (None, []) ): return None labels_value = item.get("labels") labels = labels_value if isinstance(labels_value, list) else [] repository_value = item.get("repository") repository = repository_value if isinstance(repository_value, dict) else {} return { "id": item.get("id"), "number": item.get("number"), "title": item.get("title", "") if isinstance(item.get("title"), str) else "", "body": item.get("body", "") if isinstance(item.get("body"), str) else "", "state": "open", "repository": repository.get("full_name", "") if isinstance(repository.get("full_name"), str) else "", "labels": [ label["name"] for label in labels if isinstance(label, dict) and isinstance(label.get("name"), str) ], "assignees": [], "updated_at": item.get("updated_at", "") if isinstance(item.get("updated_at"), str) else "", "url": _safe_web_url(item.get("html_url")), } async def available_issue_snapshot(max_pages: int = 10, upstream_limit: int = 50) -> list[dict]: """Load, filter, and globally rank a bounded snapshot of available issues.""" items: list[dict] = [] seen_ids: set[Any] = set() async def load_page(upstream_page: int) -> tuple[list[Any], int | None]: response = await _get_client().get( "/api/v1/repos/issues/search", headers=_auth(), params={ "state": "open", "type": "issues", "limit": upstream_limit, "page": upstream_page, }, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea available issue search response was not a list") try: total = max(0, int(response.headers["X-Total-Count"])) except (KeyError, TypeError, ValueError): total = None return payload, total first_payload, upstream_total = await load_page(1) pages: list[list[Any]] = [first_payload] if upstream_total is not None: page_count = min(max_pages, max(1, (upstream_total + upstream_limit - 1) // upstream_limit)) semaphore = asyncio.Semaphore(AVAILABLE_ISSUE_PAGE_CONCURRENCY) async def load_bounded(page: int) -> list[Any]: async with semaphore: payload, _ = await load_page(page) return payload if page_count > 1: pages.extend(await asyncio.gather(*( load_bounded(page) for page in range(2, page_count + 1) ))) else: previous = first_payload for page in range(2, max_pages + 1): if not previous or len(previous) < upstream_limit: break previous, _ = await load_page(page) pages.append(previous) for payload in pages: for raw_item in payload: item = _normalize_available_issue(raw_item) if item is not None and item["id"] not in seen_ids: seen_ids.add(item["id"]) items.append(item) priority = {"p0", "priority-high", "critical"} items.sort(key=lambda item: (item["repository"], item["number"] or 0)) items.sort(key=lambda item: item["updated_at"], reverse=True) items.sort(key=lambda item: ( 0 if any(str(label).lower() in priority for label in item["labels"]) else 1 )) return items async def available_issue_page(page: int = 1, limit: int = 50) -> dict: """Return one logical page from a bounded, globally ranked available-work scan.""" items = await available_issue_snapshot() total = len(items) start = (page - 1) * limit page_items = items[start:start + limit] return { "items": page_items, "page": page, "total": total, "has_more": start + len(page_items) < total, } def _page_metadata(result: dict) -> dict: return { "page": result["page"], "total": result["total"], "has_more": result["has_more"], } async def issues() -> WorkItems: result = await work_page("issue") return WorkItems(result["items"], {"issue": _page_metadata(result)}) def _safe_web_url(value: Any) -> str: if not isinstance(value, str): return "" parsed = urlsplit(value) return value if parsed.scheme in {"http", "https"} and parsed.netloc else "" def _safe_gitea_web_url(value: Any) -> str: if not isinstance(value, str) or not value.strip(): return "" resolved = urljoin(f"{GITEA_URL}/", value.strip()) parsed = urlsplit(resolved) configured = urlsplit(GITEA_URL) base_path = configured.path.rstrip("/") if ( parsed.scheme not in {"http", "https"} or parsed.scheme != configured.scheme or parsed.netloc != configured.netloc or (base_path and parsed.path != base_path and not parsed.path.startswith(f"{base_path}/")) ): return "" return resolved def _normalize_commit_checks(status: Any) -> list[dict]: entries = status.get("statuses") if isinstance(status, dict) else None if not isinstance(entries, list): return [] rank = {"error": 0, "failure": 0, "pending": 1, "warning": 2, "success": 3} checks = [] for index, entry in enumerate(entries): if not isinstance(entry, dict): continue name = entry.get("context") if not isinstance(name, str) or not name.strip(): continue state = entry.get("status", entry.get("state", "unknown")) state = state.lower() if isinstance(state, str) else "unknown" if state not in rank: state = "unknown" description = entry.get("description") checks.append({ "name": name.strip()[:120], "state": state, "description": description.strip()[:240] if isinstance(description, str) else "", "url": _safe_gitea_web_url(entry.get("target_url")), "_index": index, }) checks.sort(key=lambda check: (rank.get(check["state"], 2), check["_index"])) return [{key: value for key, value in check.items() if key != "_index"} for check in checks[:20]] def _normalize_notifications(threads: Any) -> list[dict]: if not isinstance(threads, list): raise ValueError("Gitea notification response was not a list") normalized = [] for thread in threads: if not isinstance(thread, dict): continue repository = thread.get("repository") subject = thread.get("subject") repository = repository if isinstance(repository, dict) else {} subject = subject if isinstance(subject, dict) else {} subject_url = _safe_web_url(subject.get("html_url")) latest_url = _safe_web_url(subject.get("latest_comment_html_url")) number_text = ( urlsplit(subject_url).path.rstrip("/").rsplit("/", 1)[-1] if subject_url else "" ) normalized.append( { "id": thread.get("id"), "unread": thread.get("unread") is True, "updated_at": ( thread.get("updated_at") if isinstance(thread.get("updated_at"), str) else "" ), "repository": ( repository.get("full_name") if isinstance(repository.get("full_name"), str) else "" ), "number": int(number_text) if number_text.isdigit() else None, "title": ( subject.get("title") if isinstance(subject.get("title"), str) and subject.get("title") else "Untitled update" ), "subject_type": ( subject.get("type") if isinstance(subject.get("type"), str) and subject.get("type") else "Update" ), "state": ( subject.get("state") if isinstance(subject.get("state"), str) else "" ), "url": latest_url or subject_url, "subject_url": subject_url, } ) return normalized async def notifications() -> dict: return await notification_page(1) async def notification_page(page: int, limit: int = 50) -> dict: response = await _get_client().get( f"/api/v1/notifications?status-types=unread&limit={limit}&page={page}", headers=_auth(), ) response.raise_for_status() items = _normalize_notifications(response.json()) try: total = max(len(items), int(response.headers.get("X-Total-Count", len(items)))) except (TypeError, ValueError): total = len(items) return { "items": items, "page": page, "total": total, "has_more": page * limit < total, } async def mark_notification_read(thread_id: int) -> None: response = await _get_client().patch( f"/api/v1/notifications/threads/{thread_id}?to-status=read", headers=_auth(), ) response.raise_for_status() async def acknowledge_notification(thread_id: int) -> dict: thread = await fetch(f"notifications/threads/{thread_id}") if not isinstance(thread, dict): raise ValueError("Gitea notification thread response was not an object") repository = thread.get("repository") subject = thread.get("subject") if not isinstance(repository, dict) or not isinstance(subject, dict): raise ValueError("Notification does not identify a conversation") repository_name = repository.get("full_name") subject_path = _gitea_api_path(subject.get("url")) comment_path = _gitea_api_path(subject.get("latest_comment_url")) subject_match = re.fullmatch( r"repos/([^/]+/[^/]+)/(issues|pulls)/(\d+)", subject_path ) comment_match = re.fullmatch( r"repos/([^/]+/[^/]+)/issues/comments/(\d+)", comment_path ) if ( not subject_match or not comment_match or subject_match.group(1) != repository_name or comment_match.group(1) != repository_name or subject.get("type") not in {"Issue", "Pull"} ): raise ValueError("Notification has no supported latest comment") user = await current_user() login = user.get("login") if isinstance(user, dict) else None if not isinstance(login, str) or not login: raise ValueError("Authenticated Gitea user is unavailable") reaction_path = f"/api/v1/{comment_path}/reactions" response = await _get_client().get(reaction_path, headers=_auth()) response.raise_for_status() reactions = response.json() existing = any( isinstance(reaction, dict) and reaction.get("content") == "+1" and isinstance(reaction.get("user"), dict) and reaction["user"].get("login") == login for reaction in (reactions if isinstance(reactions, list) else []) ) if not existing: response = await _get_client().post( reaction_path, headers=_auth(), json={"content": "+1"} ) response.raise_for_status() await mark_notification_read(thread_id) return { "id": thread_id, "reaction": "existing" if existing else "created", "status": "read", } def _gitea_api_path(value: Any) -> str: if not isinstance(value, str): return "" parsed = urlsplit(value) configured = urlsplit(GITEA_URL) prefix = "/api/v1/" if ( parsed.scheme not in {"http", "https"} or parsed.netloc != configured.netloc or not parsed.path.startswith(prefix) ): return "" return parsed.path[len(prefix):] + (("?" + parsed.query) if parsed.query else "") async def notification_detail(thread_id: int) -> dict: thread = await fetch(f"notifications/threads/{thread_id}") if not isinstance(thread, dict): raise ValueError("Gitea notification thread response was not an object") repository_value = thread.get("repository") repository = repository_value if isinstance(repository_value, dict) else {} subject_value = thread.get("subject") subject = subject_value if isinstance(subject_value, dict) else {} subject_path = _gitea_api_path(subject.get("url")) comment_path = _gitea_api_path(subject.get("latest_comment_url")) conversation_match = re.fullmatch( r"repos/([^/]+/[^/]+)/(issues|pulls)/(\d+)", subject_path ) comment_match = re.fullmatch( r"repos/([^/]+/[^/]+)/issues/comments/(\d+)", comment_path ) repository_name = repository.get("full_name") supported_conversation = ( conversation_match and conversation_match.group(1) == repository_name and subject.get("type") in {"Issue", "Pull"} ) if supported_conversation: assert conversation_match is not None subject_detail, comment, conversation = await asyncio.gather( fetch(subject_path), fetch(comment_path) if comment_path else asyncio.sleep(0, result={}), issue_conversation_page( conversation_match.group(1), int(conversation_match.group(3)) ), ) else: subject_detail, comment, conversation = await asyncio.gather( fetch(subject_path) if subject_path else asyncio.sleep(0, result={}), fetch(comment_path) if comment_path else asyncio.sleep(0, result={}), asyncio.sleep(0, result={"comments": [], "page": 1, "older_page": None, "total": 0}), ) subject_detail = subject_detail if isinstance(subject_detail, dict) else {} comment = comment if isinstance(comment, dict) else {} user_value = comment.get("user") user = user_value if isinstance(user_value, dict) else {} subject_url = _safe_web_url(subject.get("html_url")) latest_url = _safe_web_url(comment.get("html_url")) or _safe_web_url( subject.get("latest_comment_html_url") ) assignee_values = subject_detail.get("assignees") assignees = [ value.get("login") for value in (assignee_values if isinstance(assignee_values, list) else []) if isinstance(value, dict) and isinstance(value.get("login"), str) ] issue = { "number": subject_detail.get("number"), "assignees": assignees, "claimable": ( subject.get("type") == "Issue" and subject_detail.get("state", subject.get("state")) == "open" and not assignees ), } if supported_conversation and subject.get("type") == "Issue" else None return { "id": thread_id, "repository": repository.get("full_name", "") if isinstance(repository.get("full_name"), str) else "", "title": subject.get("title", "") if isinstance(subject.get("title"), str) else "", "subject_type": subject.get("type", "Update") if isinstance(subject.get("type"), str) else "Update", "state": subject.get("state", "") if isinstance(subject.get("state"), str) else "", "url": latest_url or subject_url, "subject_body": subject_detail.get("body", "") if isinstance(subject_detail.get("body"), str) else "", "latest_comment": { "author": user.get("login", "") if isinstance(user.get("login"), str) else "", "body": comment.get("body", "") if isinstance(comment.get("body"), str) else "", "created_at": comment.get("created_at", "") if isinstance(comment.get("created_at"), str) else "", "url": latest_url, }, "issue": issue, "acknowledge_supported": bool( supported_conversation and comment_match and comment_match.group(1) == repository_name ), "conversation": conversation, } async def notification_conversation_page( thread_id: int, page: int, limit: int = 20 ) -> dict: thread = await fetch(f"notifications/threads/{thread_id}") if not isinstance(thread, dict): raise ValueError("Gitea notification thread response was not an object") repository = thread.get("repository") subject = thread.get("subject") if not isinstance(repository, dict) or not isinstance(subject, dict): raise ValueError("Notification does not identify a conversation") subject_path = _gitea_api_path(subject.get("url")) match = re.fullmatch(r"repos/([^/]+/[^/]+)/(issues|pulls)/(\d+)", subject_path) if ( not match or match.group(1) != repository.get("full_name") or subject.get("type") not in {"Issue", "Pull"} ): raise ValueError("Notification subject is not a supported conversation") return await issue_conversation_page( match.group(1), int(match.group(3)), page=page, limit=limit ) async def reply_to_notification(thread_id: int, body: str) -> dict: thread = await fetch(f"notifications/threads/{thread_id}") if not isinstance(thread, dict): raise ValueError("Gitea notification thread response was not an object") repository = thread.get("repository") subject = thread.get("subject") if not isinstance(repository, dict) or not isinstance(subject, dict): raise ValueError("Notification does not identify a conversation") repository_name = repository.get("full_name") subject_path = _gitea_api_path(subject.get("url")) match = re.fullmatch( r"repos/([^/]+/[^/]+)/(issues|pulls)/(\d+)", subject_path ) if ( not match or match.group(1) != repository_name or subject.get("type") not in {"Issue", "Pull"} ): raise ValueError("Notification subject is not a supported conversation") response = await _get_client().post( f"/api/v1/repos/{match.group(1)}/issues/{match.group(3)}/comments", headers=_auth(), json={"body": body}, ) response.raise_for_status() comment = response.json() if not isinstance(comment, dict): raise ValueError("Gitea comment response was not an object") return _normalize_issue_comment(comment) async def upload_notification_attachment( thread_id: int, filename: str, content_type: str, content: bytes ) -> dict: """Upload to the exact issue or pull identified by a trusted notification.""" thread = await fetch(f"notifications/threads/{thread_id}") if not isinstance(thread, dict): raise ValueError("Gitea notification thread response was not an object") repository = thread.get("repository") subject = thread.get("subject") if not isinstance(repository, dict) or not isinstance(subject, dict): raise ValueError("Notification does not identify a conversation") subject_path = _gitea_api_path(subject.get("url")) match = re.fullmatch(r"repos/([^/]+/[^/]+)/(issues|pulls)/(\d+)", subject_path) if ( not match or match.group(1) != repository.get("full_name") or subject.get("type") not in {"Issue", "Pull"} ): raise ValueError("Notification subject is not a supported conversation") # Gitea stores pull-request assets on its shared issue asset endpoint. response = await _get_client().post( f"/api/v1/repos/{match.group(1)}/issues/{match.group(3)}/assets", headers=_auth(), params={"name": filename}, files={"attachment": (filename, content, content_type)}, ) response.raise_for_status() attachment = response.json() if not isinstance(attachment, dict): raise ValueError("Gitea attachment response was not an object") name = attachment.get("name") url = _safe_web_url(attachment.get("browser_download_url")) size = attachment.get("size") if not isinstance(name, str) or not name or not url or not isinstance(size, int): raise ValueError("Gitea did not confirm the attachment") return {"name": name, "url": url, "size": size} async def close_issue(repository: str, number: int) -> dict: response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"state": "closed"}, ) response.raise_for_status() issue = response.json() if not isinstance(issue, dict) or issue.get("state") != "closed": raise ValueError("Gitea did not confirm issue closure") return { "number": issue.get("number"), "state": "closed", "closed_at": issue.get("closed_at", "") if isinstance(issue.get("closed_at"), str) else "", } def _normalize_issue_comment(comment: dict) -> dict: user_value = comment.get("user") user: dict = user_value if isinstance(user_value, dict) else {} return { "id": comment.get("id"), "author": user.get("login", "") if isinstance(user.get("login"), str) else "", "body": comment.get("body", "") if isinstance(comment.get("body"), str) else "", "created_at": comment.get("created_at", "") if isinstance(comment.get("created_at"), str) else "", "url": _safe_web_url(comment.get("html_url")), } async def comment_on_issue(repository: str, number: int, body: str) -> dict: response = await _get_client().post( f"/api/v1/repos/{repository}/issues/{number}/comments", headers=_auth(), json={"body": body}, ) response.raise_for_status() comment = response.json() if not isinstance(comment, dict): raise ValueError("Gitea comment response was not an object") return _normalize_issue_comment(comment) async def upload_assigned_issue_attachment( repository: str, number: int, filename: str, content_type: str, content: bytes, ) -> dict: if not await is_assigned_issue(repository, number): raise IssueNotAvailableError("Assigned issue not found") response = await _get_client().post( f"/api/v1/repos/{repository}/issues/{number}/assets", headers=_auth(), params={"name": filename}, files={"attachment": (filename, content, content_type)}, ) response.raise_for_status() attachment = response.json() if not isinstance(attachment, dict): raise ValueError("Gitea attachment response was not an object") name = attachment.get("name") url = _safe_web_url(attachment.get("browser_download_url")) size = attachment.get("size") if not isinstance(name, str) or not name or not url or not isinstance(size, int): raise ValueError("Gitea did not confirm the attachment") return {"name": name, "url": url, "size": size} async def upload_assigned_pull_attachment( repository: str, number: int, filename: str, content_type: str, content: bytes, ) -> dict: if not await is_assigned_pull(repository, number): raise IssueNotAvailableError("Assigned pull request not found") response = await _get_client().post( f"/api/v1/repos/{repository}/issues/{number}/assets", headers=_auth(), params={"name": filename}, files={"attachment": (filename, content, content_type)}, ) response.raise_for_status() attachment = response.json() if not isinstance(attachment, dict): raise ValueError("Gitea attachment response was not an object") name = attachment.get("name") url = _safe_web_url(attachment.get("browser_download_url")) size = attachment.get("size") if not isinstance(name, str) or not name or not url or not isinstance(size, int): raise ValueError("Gitea did not confirm the attachment") return {"name": name, "url": url, "size": size} async def repo_labels(repository: str) -> list[dict]: response = await _get_client().get( f"/api/v1/repos/{repository}/labels", headers=_auth(), params={"limit": 50, "page": 1}, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea labels response was not a list") return [ { "id": item["id"], "name": item["name"], "color": item.get("color", "") if isinstance(item.get("color"), str) else "", "description": item.get("description", "") if isinstance(item.get("description"), str) else "", } for item in payload if isinstance(item, dict) and isinstance(item.get("id"), int) and isinstance(item.get("name"), str) ] async def repo_milestones(repository: str) -> list[dict]: response = await _get_client().get( f"/api/v1/repos/{repository}/milestones", headers=_auth(), params={"state": "open", "limit": 50, "page": 1}, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea milestones response was not a list") return [ {"id": item["id"], "title": item["title"]} for item in payload if isinstance(item, dict) and isinstance(item.get("id"), int) and isinstance(item.get("title"), str) and item.get("state") == "open" ] async def create_issue( repository: str, title: str, body: str, assignee: str, label_ids: list[int] | None = None, milestone_id: int | None = None, due_date: str | None = None, ) -> dict: payload: dict = {"title": title, "body": body, "assignee": assignee} if label_ids: payload["labels"] = label_ids if milestone_id is not None: payload["milestone"] = milestone_id if due_date is not None: payload["due_date"] = due_date response = await _get_client().post( f"/api/v1/repos/{repository}/issues", headers=_auth(), json=payload, ) response.raise_for_status() issue = response.json() if not isinstance(issue, dict) or not isinstance(issue.get("number"), int): raise ValueError("Gitea did not confirm issue creation") assignees_value = issue.get("assignees") assignees = assignees_value if isinstance(assignees_value, list) else [] confirmed_assignees = [ item["login"] for item in assignees if isinstance(item, dict) and isinstance(item.get("login"), str) ] if assignee not in confirmed_assignees: raise ValueError("Gitea did not confirm issue self-assignment") labels_value = issue.get("labels") labels = labels_value if isinstance(labels_value, list) else [] milestone_value = issue.get("milestone") milestone = ( {"id": milestone_value["id"], "title": milestone_value["title"]} if isinstance(milestone_value, dict) and isinstance(milestone_value.get("id"), int) and isinstance(milestone_value.get("title"), str) else None ) confirmed_due_date = ( issue.get("due_date") if isinstance(issue.get("due_date"), str) else None ) if ( (milestone_id is not None and (milestone or {}).get("id") != milestone_id) or (due_date is not None and confirmed_due_date != due_date) ): raise ValueError("Gitea did not confirm issue release plan") return { "id": issue.get("id"), "number": issue["number"], "title": issue.get("title", "") if isinstance(issue.get("title"), str) else "", "state": issue.get("state", "") if isinstance(issue.get("state"), str) else "", "repository": repository, "labels": [ item["name"] for item in labels if isinstance(item, dict) and isinstance(item.get("name"), str) ], "assignees": confirmed_assignees, "milestone": milestone, "due_date": confirmed_due_date, "updated_at": issue.get("updated_at", "") if isinstance(issue.get("updated_at"), str) else "", "url": _safe_web_url(issue.get("html_url")), } async def claim_available_issue(repository: str, number: int) -> dict: issue = await fetch(f"repos/{repository}/issues/{number}") user = await current_user() assignees_value = issue.get("assignees") if isinstance(issue, dict) else None if ( not isinstance(issue, dict) or issue.get("state") != "open" or issue.get("pull_request") is not None or assignees_value not in (None, []) ): raise IssueNotAvailableError("Issue is no longer available") login = user.get("login") if isinstance(user, dict) else None if not isinstance(login, str) or not login: raise ValueError("Gitea current user did not include a login") response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"assignee": login}, ) response.raise_for_status() confirmed = response.json() if not isinstance(confirmed, dict) or confirmed.get("number") != number: raise ValueError("Gitea did not confirm issue assignment") confirmed_assignees_value = confirmed.get("assignees") confirmed_assignees = ( confirmed_assignees_value if isinstance(confirmed_assignees_value, list) else [] ) logins = [ assignee["login"] for assignee in confirmed_assignees if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) ] if login not in logins: raise ValueError("Gitea did not confirm issue assignment") labels_value = confirmed.get("labels") labels = labels_value if isinstance(labels_value, list) else [] return { "id": confirmed.get("id"), "number": number, "title": confirmed.get("title", "") if isinstance(confirmed.get("title"), str) else "", "state": "open", "repository": repository, "labels": [ label["name"] for label in labels if isinstance(label, dict) and isinstance(label.get("name"), str) ], "assignees": logins, "updated_at": confirmed.get("updated_at", "") if isinstance(confirmed.get("updated_at"), str) else "", "url": _safe_web_url(confirmed.get("html_url")), } async def reopen_issue(repository: str, number: int) -> dict: """Reopen a closed issue, assign it to the current user, and confirm both.""" issue, user = await asyncio.gather( fetch(f"repos/{repository}/issues/{number}"), current_user() ) login = user.get("login") if isinstance(user, dict) else None if ( not isinstance(issue, dict) or issue.get("number") != number or issue.get("pull_request") is not None or not isinstance(login, str) or not login ): raise IssueNotAvailableError("Issue cannot be resumed") if issue.get("state") == "open" and _login_in_users(login, issue.get("assignees")): confirmed = issue elif issue.get("state") == "closed": response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"state": "open", "assignee": login}, ) response.raise_for_status() confirmed = response.json() else: raise IssueNotAvailableError("Issue is no longer available to resume") assignees_value = confirmed.get("assignees") if isinstance(confirmed, dict) else None assignees = assignees_value if isinstance(assignees_value, list) else [] logins = [ assignee["login"] for assignee in assignees if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) ] if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or confirmed.get("state") != "open" or login not in logins ): raise ValueError("Gitea did not confirm issue reopening and assignment") labels_value = confirmed.get("labels") labels = labels_value if isinstance(labels_value, list) else [] return { "id": confirmed.get("id"), "number": number, "title": confirmed.get("title", "") if isinstance(confirmed.get("title"), str) else "", "state": "open", "repository": repository, "labels": [ label["name"] for label in labels if isinstance(label, dict) and isinstance(label.get("name"), str) ], "assignees": logins, "updated_at": confirmed.get("updated_at", "") if isinstance(confirmed.get("updated_at"), str) else "", "url": _safe_web_url(confirmed.get("html_url")), } async def release_assigned_issue(repository: str, number: int) -> dict: login, issue = await _current_login_and_target( f"repos/{repository}/issues/{number}" ) assignees_value = issue.get("assignees") if ( issue.get("state") != "open" or issue.get("pull_request") is not None or not _login_in_users(login, assignees_value) ): raise IssueNotAvailableError("Issue is not assigned to the current user") assignees = assignees_value if isinstance(assignees_value, list) else [] remaining = [ assignee["login"] for assignee in assignees if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) and assignee["login"] != login ] response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"assignees": remaining}, ) response.raise_for_status() confirmed = response.json() confirmed_value = confirmed.get("assignees") if isinstance(confirmed, dict) else None confirmed_assignees = [ assignee["login"] for assignee in confirmed_value if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) ] if isinstance(confirmed_value, list) else [] if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or login in confirmed_assignees or set(confirmed_assignees) != set(remaining) ): raise ValueError("Gitea did not confirm issue release") return { "number": number, "repository": repository, "state": confirmed.get("state", "open"), "assignees": confirmed_assignees, "available": not confirmed_assignees, } async def issue_handoff_candidates(repository: str) -> list[dict]: user, response = await asyncio.gather( current_user(), _get_client().get( f"/api/v1/repos/{repository}/assignees", headers=_auth() ), ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea assignees response was not a list") current_login = user.get("login") if isinstance(user, dict) else None candidates = [] for item in payload: if not isinstance(item, dict): continue login = item.get("login") if not isinstance(login, str) or not login or login == current_login: continue full_name = item.get("full_name") candidates.append({ "login": login, "name": full_name if isinstance(full_name, str) and full_name else login, }) return candidates async def mention_candidates( repository: str, query: str, *, limit: int = 8 ) -> list[dict]: response = await _get_client().get( f"/api/v1/repos/{repository}/assignees", headers=_auth() ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea assignees response was not a list") needle = query.casefold() matches = [] for item in payload: if not isinstance(item, dict): continue login = item.get("login") if ( not isinstance(login, str) or not re.fullmatch(r"[A-Za-z0-9_.-]+", login) ): continue full_name = item.get("full_name") name = full_name if isinstance(full_name, str) and full_name else login login_match = needle in login.casefold() name_match = needle in name.casefold() if login_match or name_match: matches.append((0 if login.casefold().startswith(needle) else 1, login.casefold(), { "login": login, "name": name, })) matches.sort(key=lambda item: (item[0], item[1])) return [item[2] for item in matches[:limit]] async def handoff_assigned_issue( repository: str, number: int, recipient: str ) -> dict: login, issue = await _current_login_and_target( f"repos/{repository}/issues/{number}" ) if ( issue.get("state") != "open" or issue.get("pull_request") is not None or not _login_in_users(login, issue.get("assignees")) ): raise IssueNotAvailableError("Issue is not assigned to the current user") eligible = { item["login"] for item in await issue_handoff_candidates(repository) } if recipient not in eligible: raise IssueNotAvailableError("Handoff recipient is not eligible") assignees_value = issue.get("assignees") assignees = assignees_value if isinstance(assignees_value, list) else [] desired = [ item["login"] for item in assignees if isinstance(item, dict) and isinstance(item.get("login"), str) and item["login"] != login ] if recipient not in desired: desired.append(recipient) response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"assignees": desired}, ) response.raise_for_status() confirmed = response.json() confirmed_value = confirmed.get("assignees") if isinstance(confirmed, dict) else None confirmed_assignees = [ item["login"] for item in confirmed_value if isinstance(item, dict) and isinstance(item.get("login"), str) ] if isinstance(confirmed_value, list) else [] if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or login in confirmed_assignees or recipient not in confirmed_assignees or set(confirmed_assignees) != set(desired) ): raise ValueError("Gitea did not confirm issue handoff") return { "repository": repository, "number": number, "state": confirmed.get("state", "open"), "assignees": confirmed_assignees, "recipient": recipient, } async def update_issue_labels(repository: str, number: int, label_ids: list[int]) -> dict: response = await _get_client().patch( f"/api/v1/repos/{repository}/issues/{number}", headers=_auth(), json={"labels": label_ids}, ) response.raise_for_status() issue = response.json() if not isinstance(issue, dict) or issue.get("number") != number: raise ValueError("Gitea did not confirm the label update") labels_value = issue.get("labels") labels = labels_value if isinstance(labels_value, list) else [] confirmed_ids = { item["id"] for item in labels if isinstance(item, dict) and isinstance(item.get("id"), int) } if confirmed_ids != set(label_ids): raise ValueError("Gitea did not confirm the requested label set") return { "number": number, "labels": [ item["name"] for item in labels if isinstance(item, dict) and isinstance(item.get("name"), str) ], } async def update_assigned_issue_due_date( repository: str, number: int, due_date: str | None ) -> dict: path = f"repos/{repository}/issues/{number}" login, issue = await _current_login_and_target(path) if ( issue.get("state") != "open" or isinstance(issue.get("pull_request"), dict) or not _login_in_users(login, issue.get("assignees")) ): raise IssueNotAvailableError("assigned issue not found") payload = {"due_date": due_date} if due_date else {"unset_due_date": True} response = await _get_client().patch( f"/api/v1/{path}", headers=_auth(), json=payload ) response.raise_for_status() confirmed = response.json() confirmed_due_date = confirmed.get("due_date") if isinstance(confirmed, dict) else None if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or confirmed_due_date != due_date ): raise ValueError("Gitea did not confirm the issue due date update") return { "repository": repository, "number": number, "state": confirmed.get("state", "open"), "due_date": confirmed_due_date, } async def update_assigned_issue_milestone( repository: str, number: int, milestone_id: int | None ) -> dict: path = f"repos/{repository}/issues/{number}" login, issue = await _current_login_and_target(path) if ( issue.get("state") != "open" or isinstance(issue.get("pull_request"), dict) or not _login_in_users(login, issue.get("assignees")) ): raise IssueNotAvailableError("assigned issue not found") selected = None if milestone_id is not None: selected = next( (item for item in await repo_milestones(repository) if item["id"] == milestone_id), None, ) if selected is None: raise ValueError("Unknown open repository milestone") response = await _get_client().patch( f"/api/v1/{path}", headers=_auth(), json={"milestone": milestone_id or 0} ) response.raise_for_status() confirmed = response.json() milestone_value = confirmed.get("milestone") if isinstance(confirmed, dict) else None normalized = ( {"id": milestone_value["id"], "title": milestone_value["title"]} if isinstance(milestone_value, dict) and isinstance(milestone_value.get("id"), int) and isinstance(milestone_value.get("title"), str) else None ) if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or normalized != selected ): raise ValueError("Gitea did not confirm the issue milestone update") return { "repository": repository, "number": number, "state": confirmed.get("state", "open"), "milestone": normalized, } async def issue_conversation_page( repository: str, number: int, page: int | None = None, limit: int = 20, ) -> dict: """Return one bounded comment page, opening on the newest page by default.""" bounded_limit = min(50, max(1, limit)) requested_page = max(1, page or 1) path = f"/api/v1/repos/{repository}/issues/{number}/comments" async def load(selected_page: int) -> tuple[list, int]: response = await _get_client().get( path, headers=_auth(), params={"limit": bounded_limit, "page": selected_page}, ) response.raise_for_status() payload = response.json() if not isinstance(payload, list): raise ValueError("Gitea issue comments response was not a list") comments = [item for item in payload if isinstance(item, dict)] try: total = max(len(comments), int(response.headers.get("X-Total-Count", len(comments)))) except (TypeError, ValueError): total = len(comments) return comments, total comments, total = await load(requested_page) newest_page = max(1, (total + bounded_limit - 1) // bounded_limit) if len(comments) >= total: selected_page = newest_page if page is None else min(requested_page, newest_page) start = (selected_page - 1) * bounded_limit comments = comments[start:start + bounded_limit] else: selected_page = requested_page if page is None and newest_page != requested_page: selected_page = newest_page comments, confirmed_total = await load(selected_page) total = max(total, confirmed_total) return { "comments": [_normalize_issue_comment(item) for item in comments], "page": selected_page, "older_page": selected_page - 1 if selected_page > 1 else None, "total": total, } async def issue_detail(repository: str, number: int) -> dict: base = f"repos/{repository}/issues/{number}" async def load_dependencies() -> tuple[bool, list[dict]]: try: return True, await issue_dependencies(repository, number) except Exception: return False, [] issue, conversation, dependency_result = await asyncio.gather( fetch(base), issue_conversation_page(repository, number), load_dependencies(), ) dependencies_available, dependencies = dependency_result if not isinstance(issue, dict): raise ValueError("Gitea issue response was not an object") labels_value = issue.get("labels") labels: list = labels_value if isinstance(labels_value, list) else [] assignees_value = issue.get("assignees") assignees: list = assignees_value if isinstance(assignees_value, list) else [] normalized_comments = conversation["comments"] return { "repository": repository, "number": number, "title": issue.get("title", "") if isinstance(issue.get("title"), str) else "", "state": issue.get("state", "") if isinstance(issue.get("state"), str) else "", "body": issue.get("body", "") if isinstance(issue.get("body"), str) else "", "updated_at": issue.get("updated_at", "") if isinstance(issue.get("updated_at"), str) else "", "due_date": issue.get("due_date") if isinstance(issue.get("due_date"), str) else None, "milestone": ( {"id": issue["milestone"]["id"], "title": issue["milestone"]["title"]} if isinstance(issue.get("milestone"), dict) and isinstance(issue["milestone"].get("id"), int) and isinstance(issue["milestone"].get("title"), str) else None ), "url": _safe_web_url(issue.get("html_url")), "labels": [ label["name"] for label in labels if isinstance(label, dict) and isinstance(label.get("name"), str) ], "assignees": [ assignee["login"] for assignee in assignees if isinstance(assignee, dict) and isinstance(assignee.get("login"), str) ], "dependencies_available": dependencies_available, "dependencies": dependencies, "comments": normalized_comments, "conversation": conversation, } async def issue_dependencies(repository: str, number: int, limit: int = 20) -> list[dict]: """Return bounded open prerequisites for an issue.""" response = await _get_client().get( f"/api/v1/repos/{repository}/issues/{number}/dependencies", headers=_auth(), params={"limit": limit}, ) response.raise_for_status() value = response.json() items = value if isinstance(value, list) else [] dependencies = [] for item in items[:limit]: if not isinstance(item, dict) or item.get("state") != "open": continue repo_value = item.get("repository") repo = repo_value if isinstance(repo_value, dict) else {} dependency_repository = repo.get("full_name") dependency_number = item.get("number") if not isinstance(dependency_repository, str) or not isinstance(dependency_number, int): continue dependencies.append({ "repository": dependency_repository, "number": dependency_number, "title": item.get("title", "") if isinstance(item.get("title"), str) else "", "state": "open", "url": _safe_web_url(item.get("html_url")), }) return dependencies async def mutate_assigned_issue_dependency( repository: str, number: int, blocker_repository: str, blocker_number: int, remove: bool = False, ) -> dict: """Add or remove one prerequisite and return the canonical open dependency set.""" if (repository, number) == (blocker_repository, blocker_number): raise IssueDependencyInvalidError("an issue cannot block itself") login, source = await _current_login_and_target( f"repos/{repository}/issues/{number}" ) if ( source.get("state") != "open" or isinstance(source.get("pull_request"), dict) or not _login_in_users(login, source.get("assignees")) ): raise IssueNotAvailableError("assigned issue not found") current = await issue_dependencies(repository, number) currently_present = any( item["repository"] == blocker_repository and item["number"] == blocker_number for item in current ) if currently_present != remove: return { "repository": repository, "number": number, "dependencies_available": True, "dependencies": current, } if not remove: candidate = await fetch(f"repos/{blocker_repository}/issues/{blocker_number}") if ( not isinstance(candidate, dict) or candidate.get("number") != blocker_number or candidate.get("state") != "open" or isinstance(candidate.get("pull_request"), dict) ): raise IssueDependencyInvalidError("blocker must be an accessible open issue") owner, repo = blocker_repository.split("/", 1) response = await _get_client().request( "DELETE" if remove else "POST", f"/api/v1/repos/{repository}/issues/{number}/dependencies", headers=_auth(), json={"owner": owner, "repo": repo, "index": blocker_number}, ) response.raise_for_status() confirmed = await issue_dependencies(repository, number) present = any( item["repository"] == blocker_repository and item["number"] == blocker_number for item in confirmed ) if present == remove: raise ValueError("Gitea did not confirm the dependency change") return { "repository": repository, "number": number, "dependencies_available": True, "dependencies": confirmed, } async def update_assigned_issue( repository: str, number: int, title: str, body: str, expected_updated_at: str, ) -> dict: path = f"repos/{repository}/issues/{number}" login, issue = await _current_login_and_target(path) if ( issue.get("state") != "open" or isinstance(issue.get("pull_request"), dict) or not _login_in_users(login, issue.get("assignees")) ): raise IssueNotAvailableError("assigned issue not found") if issue.get("updated_at") != expected_updated_at: raise IssueEditConflictError("issue changed upstream") response = await _get_client().patch( f"/api/v1/{path}", headers=_auth(), json={"title": title, "body": body}, ) response.raise_for_status() confirmed = response.json() if ( not isinstance(confirmed, dict) or confirmed.get("number") != number or confirmed.get("title") != title or confirmed.get("body", "") != body ): raise ValueError("Gitea did not confirm the issue content update") return { "repository": repository, "number": number, "title": title, "body": body, "state": confirmed.get("state", "open"), "updated_at": confirmed.get("updated_at", ""), "url": _safe_web_url(confirmed.get("html_url")), } async def _current_login_and_target(path: str) -> tuple[str, dict]: user, target = await asyncio.gather(current_user(), fetch(path)) login = user.get("login") if isinstance(user, dict) else None if not isinstance(login, str) or not login or not isinstance(target, dict): return "", {} return login, target def _login_in_users(login: str, value: object) -> bool: return isinstance(value, list) and any( isinstance(user, dict) and user.get("login") == login for user in value ) async def resolve_work_route( kind: str, repository: str | None, number: int | None, notification_id: int | None, ) -> dict: if kind == "update": if notification_id is None: raise WorkRouteUnavailableError("Notification identity is missing") thread = await fetch(f"notifications/threads/{notification_id}") if not isinstance(thread, dict) or thread.get("unread") is not True: raise WorkRouteUnavailableError("Notification is no longer unread") detail = await notification_detail(notification_id) return { "kind": "update", "notification_id": notification_id, "has_update": True, "repository": detail.get("repository", ""), "title": detail.get("title", ""), "url": detail.get("url", ""), } if repository is None or number is None or kind not in {"issue", "pull", "review"}: raise WorkRouteUnavailableError("Work identity is missing") target_path = "issues" if kind == "issue" else "pulls" login, target = await _current_login_and_target( f"repos/{repository}/{target_path}/{number}" ) assigned = _login_in_users(login, target.get("assignees")) requested = _login_in_users(login, target.get("requested_reviewers")) eligible = ( target.get("state") == "open" and ( (kind == "issue" and not isinstance(target.get("pull_request"), dict) and assigned) or (kind == "pull" and assigned) or (kind == "review" and requested) ) ) if not eligible: raise WorkRouteUnavailableError("Work item is no longer in My Work") return { "kind": kind, "repository": repository, "number": number, "title": target.get("title", "") if isinstance(target.get("title"), str) else "", "state": "open", "url": _safe_web_url(target.get("html_url")), **({"is_review": True, "work_reasons": ["review_requested"]} if kind == "review" else {}), } async def is_assigned_issue(repository: str, number: int) -> bool: login, issue = await _current_login_and_target( f"repos/{repository}/issues/{number}" ) return ( issue.get("state") == "open" and not isinstance(issue.get("pull_request"), dict) and _login_in_users(login, issue.get("assignees")) ) async def pull_requests() -> WorkItems: assigned, review_requested = await asyncio.gather( work_page("pull"), work_page("review"), ) merged: dict[int, dict] = {} for result in ( assigned, review_requested, ): for pull in result["items"]: identity = pull.get("id") if identity not in merged: merged[identity] = {**pull, "work_reasons": []} for reason in pull.get("work_reasons", []): if reason not in merged[identity]["work_reasons"]: merged[identity]["work_reasons"].append(reason) return WorkItems( list(merged.values()), { "pull": _page_metadata(assigned), "review": _page_metadata(review_requested), }, ) async def is_requested_review(repository: str, number: int) -> bool: login, pull = await _current_login_and_target( f"repos/{repository}/pulls/{number}" ) return ( pull.get("state") == "open" and _login_in_users(login, pull.get("requested_reviewers")) ) async def is_assigned_pull(repository: str, number: int) -> bool: login, pull = await _current_login_and_target( f"repos/{repository}/pulls/{number}" ) return ( pull.get("state") == "open" and _login_in_users(login, pull.get("assignees")) ) async def pull_completion_detail(repository: str, number: int) -> dict: base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) if not isinstance(pull, dict): raise ValueError("Gitea pull request response was not an object") conversation = await issue_conversation_page(repository, number) head = pull.get("head") if isinstance(pull.get("head"), dict) else {} sha = head.get("sha") if isinstance(head.get("sha"), str) else "" user = pull.get("user") if isinstance(pull.get("user"), dict) else {} return { "repository": repository, "number": number, "title": pull.get("title") if isinstance(pull.get("title"), str) else "", "body": pull.get("body") if isinstance(pull.get("body"), str) else "", "url": _safe_web_url(pull.get("html_url")), "author": user.get("login") if isinstance(user.get("login"), str) else "", "head_sha": sha, "state": pull.get("state") if isinstance(pull.get("state"), str) else "", "conversation": conversation, } async def pull_completion_review(repository: str, number: int) -> dict: base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) if not isinstance(pull, dict): raise ValueError("Gitea pull request response was not an object") head_value = pull.get("head") head: dict = head_value if isinstance(head_value, dict) else {} sha_value = head.get("sha") sha = sha_value if isinstance(sha_value, str) else "" files, status, diff_result = await asyncio.gather( fetch(f"{base}/files"), fetch(f"repos/{repository}/commits/{sha}/status"), fetch_text( f"repos/{repository}/pulls/{number}.diff", REVIEW_DIFF_MAX_BYTES ), ) diff, diff_truncated = diff_result previews = _diff_previews(diff, diff_truncated) return { "repository": repository, "number": number, "head_sha": sha, "state": pull.get("state") if isinstance(pull.get("state"), str) else "", "draft": pull.get("draft") is True, "mergeable": pull.get("mergeable") is True, "merged": pull.get("merged") is True, "ci_state": status.get("state", "unknown") if isinstance(status, dict) else "unknown", "checks": _normalize_commit_checks(status), "files": [ { "filename": item.get("filename", ""), "status": item.get("status") or "changed", "additions": item.get("additions") or 0, "deletions": item.get("deletions") or 0, **previews.get( item["filename"], { "diff_lines": [], "diff_available": False, "diff_binary": False, "diff_truncated": diff_truncated, }, ), } for item in (files if isinstance(files, list) else [])[:100] if isinstance(item, dict) and isinstance(item.get("filename"), str) ], } async def pull_check_status(repository: str, number: int) -> dict: """Load only mutable pull and CI state, without immutable review data.""" base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) if not isinstance(pull, dict): raise ValueError("Gitea pull request response was not an object") head_value = pull.get("head") head: dict = head_value if isinstance(head_value, dict) else {} sha_value = head.get("sha") sha = sha_value if isinstance(sha_value, str) else "" status = await fetch(f"repos/{repository}/commits/{sha}/status") return { "repository": repository, "number": number, "head_sha": sha, "state": pull.get("state") if isinstance(pull.get("state"), str) else "", "draft": pull.get("draft") is True, "mergeable": pull.get("mergeable") is True, "merged": pull.get("merged") is True, "ci_state": status.get("state", "unknown") if isinstance(status, dict) else "unknown", "checks": _normalize_commit_checks(status), } async def is_pull_merged_at_head( repository: str, number: int, expected_head_sha: str ) -> bool: pull = await fetch(f"repos/{repository}/pulls/{number}") if not isinstance(pull, dict): return False head_value = pull.get("head") head = head_value if isinstance(head_value, dict) else {} return ( pull.get("merged") is True and pull.get("state") == "closed" and head.get("sha") == expected_head_sha ) async def merge_assigned_pull( repository: str, number: int, expected_head_sha: str ) -> dict: base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) if not isinstance(pull, dict): raise PullNotMergeableError("Pull request state is unavailable") head_value = pull.get("head") head: dict = head_value if isinstance(head_value, dict) else {} current_sha = head.get("sha") if current_sha != expected_head_sha: raise StalePullError("Pull request changed before merge") status = await fetch(f"repos/{repository}/commits/{current_sha}/status") ci_state = status.get("state") if isinstance(status, dict) else "unknown" if ( pull.get("state") != "open" or pull.get("draft") is True or pull.get("mergeable") is not True or pull.get("merged") is True or ci_state != "success" ): raise PullNotMergeableError("Pull request is not currently safe to merge") response = await _get_client().post( f"/api/v1/{base}/merge", headers=_auth(), json={"Do": "merge", "head_commit_id": current_sha}, ) response.raise_for_status() return {"number": number, "merged": True, "state": "closed"} async def pull_review_detail(repository: str, number: int) -> dict: base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) if not isinstance(pull, dict): raise ValueError("Gitea pull request response was not an object") head_value = pull.get("head") head: dict = head_value if isinstance(head_value, dict) else {} sha_value = head.get("sha") sha = sha_value if isinstance(sha_value, str) else "" files, status, reviews, diff_result = await asyncio.gather( fetch(f"{base}/files"), fetch(f"repos/{repository}/commits/{sha}/status"), fetch(f"{base}/reviews"), fetch_text( f"repos/{repository}/pulls/{number}.diff", REVIEW_DIFF_MAX_BYTES ), ) diff, diff_truncated = diff_result previews = _diff_previews(diff, diff_truncated) user_value = pull.get("user") user: dict = user_value if isinstance(user_value, dict) else {} normalized_files = [ { "filename": file["filename"], "status": file.get("status") or "changed", "additions": file.get("additions") or 0, "deletions": file.get("deletions") or 0, **previews.get( file["filename"], { "diff_lines": [], "diff_available": False, "diff_binary": False, "diff_truncated": diff_truncated, }, ), } for file in (files if isinstance(files, list) else [])[:100] if isinstance(file, dict) and isinstance(file.get("filename"), str) ] normalized_reviews = [ { "user": { "login": ( review["user"].get("login", "") if isinstance(review.get("user"), dict) else "" ) }, "state": review.get("state") or "COMMENT", "body": review.get("body") or "", } for review in (reviews if isinstance(reviews, list) else [])[:50] if isinstance(review, dict) ] return { "repository": repository, "number": number, "title": pull.get("title", ""), "body": pull.get("body") or "", "url": pull.get("html_url", ""), "author": user.get("login", ""), "head_sha": sha, "ci_state": status.get("state", "unknown") if isinstance(status, dict) else "unknown", "checks": _normalize_commit_checks(status), "files": normalized_files, "reviews": normalized_reviews, } async def submit_pull_review( repository: str, number: int, expected_head_sha: str, decision: str, body: str, comments: list[dict] | None = None, ) -> dict: base = f"repos/{repository}/pulls/{number}" pull = await fetch(base) head = pull.get("head") if isinstance(pull, dict) else None current_sha = head.get("sha") if isinstance(head, dict) else None if current_sha != expected_head_sha: raise StaleReviewError("Pull request changed while it was being reviewed") if comments: files = await fetch(f"{base}/files") changed_paths = { item.get("filename") for item in (files if isinstance(files, list) else []) if isinstance(item, dict) and isinstance(item.get("filename"), str) } if any(comment.get("path") not in changed_paths for comment in comments): raise InvalidReviewCommentError("Inline comment path is not in this pull request") response = await _get_client().post( f"/api/v1/{base}/reviews", headers=_auth(), json={ "body": body, "event": { "comment": "COMMENT", "approve": "APPROVE", "request_changes": "REQUEST_CHANGES", }[decision], "commit_id": current_sha, **({"comments": comments} if comments else {}), }, ) response.raise_for_status() review = response.json() if not isinstance(review, dict): raise ValueError("Gitea review response was not an object") return { "id": review.get("id"), "state": review.get("state") or "COMMENT", "url": _safe_web_url(review.get("html_url")), } async def activity_events(user: dict | None = None) -> list[dict]: if user is None: user = await current_user() events = await fetch(f"users/{user['login']}/activities/feeds?limit=20") if events is None: events = [] elif not isinstance(events, list): raise ValueError("Gitea activity feed response was not a list") return [ { "type": ( event.get("op_type") if isinstance(event.get("op_type"), str) and event.get("op_type") else "activity" ), "actor": ( event.get("act_user") if isinstance(event.get("act_user"), dict) else {} ), "repo": ( event.get("repo") if isinstance(event.get("repo"), dict) else {} ), "created_at": ( event.get("created") if isinstance(event.get("created"), str) else "" ), } for event in (events or []) if isinstance(event, dict) ]