"""FlowDeck — v5.2.0 Projects: normalized forge-agnostic project registry. The ``projects`` table stores one row per repository across forges (builtin / gitea / github). A background scheduler refreshes metadata (default branch, language) periodically so the UI always shows up-to-date info. """ from __future__ import annotations import logging from app.config import settings from app.db import get_conn from app.services.forge_adapter import GiteaAdapter, normalize_repo logger = logging.getLogger(__name__) def register_repo(repo: dict, proj_type: str) -> int | None: """Upsert a forge repo into the projects table. Returns project id.""" data = normalize_repo(repo, proj_type) if not data["name"]: return None with get_conn() as conn: existing = conn.execute( "SELECT id FROM projects WHERE proj_type=? AND owner=? AND name=?", (proj_type, data["owner"], data["name"]), ).fetchone() if existing: conn.execute( """UPDATE projects SET forge_id=?, clone_url=?, default_branch=?, language=?, description=?, last_synced_at=CURRENT_TIMESTAMP WHERE id=?""", (data["forge_id"], data["clone_url"], data["default_branch"], data["language"], data["description"], existing["id"]), ) conn.commit() return existing["id"] cur = conn.execute( """INSERT INTO projects (name, proj_type, owner, forge_id, clone_url, default_branch, language, description, last_synced_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)""", (data["name"], proj_type, data["owner"], data["forge_id"], data["clone_url"], data["default_branch"], data["language"], data["description"]), ) conn.commit() return cur.lastrowid def list_projects(proj_type: str | None = None) -> list[dict]: with get_conn() as conn: if proj_type: rows = conn.execute( "SELECT * FROM projects WHERE proj_type=? ORDER BY name", (proj_type,), ).fetchall() else: rows = conn.execute("SELECT * FROM projects ORDER BY proj_type, name").fetchall() return [dict(r) for r in rows] def create_builtin_project(name: str, owner: str = "", description: str = "") -> dict: """Register a standalone (non-forge) project.""" with get_conn() as conn: existing = conn.execute( "SELECT id FROM projects WHERE proj_type='builtin' AND owner=? AND name=?", (owner, name), ).fetchone() if existing: return {"id": existing["id"], "name": name} cur = conn.execute( "INSERT INTO projects (name, proj_type, owner, description) VALUES (?, 'builtin', ?, ?)", (name, owner, description), ) conn.commit() return {"id": cur.lastrowid, "name": name} # ═══════════ Periodic sync ═══════════ async def _sync_gitea(user_id: int, token: str) -> int: from app.services.gitea_client import GiteaClient gitea = GiteaClient(user_token=token) adapter = GiteaAdapter(gitea) repos = await adapter.list_repos(page=1) count = 0 for repo in repos: if register_repo(repo, "gitea"): count += 1 return count async def _sync_github(user_id: int, token: str) -> int: from app.services.github_adapter import GitHubAdapter adapter = GitHubAdapter(token) repos = await adapter.list_all_repos() count = 0 for repo in repos: if register_repo(repo, "github"): count += 1 return count async def sync_all_projects() -> dict: """Refresh the projects table from every connected forge token.""" stats = {"gitea": 0, "github": 0, "error": 0} with get_conn() as conn: rows = conn.execute( "SELECT user_id, provider, access_token FROM user_oauth_tokens " "WHERE provider IN ('gitea','github') AND access_token != ''" ).fetchall() tokens = [dict(r) for r in rows] for t in tokens: try: if t["provider"] == "gitea": stats["gitea"] += await _sync_gitea(t["user_id"], t["access_token"]) elif t["provider"] == "github": stats["github"] += await _sync_github(t["user_id"], t["access_token"]) except Exception as exc: stats["error"] += 1 logger.warning("project sync (%s user=%s) failed: %s", t["provider"], t["user_id"], exc) logger.info("projects sync done: %s", stats) return stats async def project_sync_scheduler(): """Background loop: refresh projects every interval (default hourly).""" while True: if settings.project_sync_enabled: try: await sync_all_projects() except Exception as exc: logger.warning("project_sync_scheduler error: %s", exc) await __import__("asyncio").sleep(settings.project_sync_interval_hours * 3600)