From e0c36035bcbbd32ef42c8563c97316c22a000c35 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:12:25 +0200 Subject: [PATCH 1/6] fournos-ui: lint --- fournos-ui/app/db.py | 33 ++++--- fournos-ui/app/forge_discovery.py | 12 ++- fournos-ui/app/k8s_client.py | 96 ++++++++++++-------- fournos-ui/app/main.py | 142 ++++++++++++++++++------------ fournos-ui/app/watcher.py | 55 +++++++----- 5 files changed, 214 insertions(+), 124 deletions(-) diff --git a/fournos-ui/app/db.py b/fournos-ui/app/db.py index b009fe4..ff5b29a 100644 --- a/fournos-ui/app/db.py +++ b/fournos-ui/app/db.py @@ -3,8 +3,9 @@ from __future__ import annotations import logging -from datetime import datetime, timezone -from typing import Any, Sequence +from collections.abc import Sequence +from datetime import UTC, datetime +from typing import Any from uuid import uuid4 from sqlalchemy import ( @@ -17,7 +18,8 @@ func, select, ) -from sqlalchemy.dialects.postgresql import ARRAY, JSONB, insert as pg_insert +from sqlalchemy.dialects.postgresql import ARRAY, JSONB +from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine from sqlalchemy.orm import DeclarativeBase, relationship @@ -34,6 +36,7 @@ class Base(DeclarativeBase): # ORM Models # --------------------------------------------------------------------------- + class Job(Base): __tablename__ = "jobs" @@ -46,7 +49,7 @@ class Job(Base): owner = Column(String, default="", index=True) status = Column(String, default="Pending", index=True) message = Column(Text, default="") - created_at = Column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc)) + created_at = Column(DateTime(timezone=True), default=lambda: datetime.now(UTC)) completed_at = Column(DateTime(timezone=True), nullable=True) duration_seconds = Column(Float, nullable=True) mlflow_url = Column(String, default="") @@ -59,17 +62,21 @@ class Job(Base): triggered_by_schedule = Column(String, nullable=True, index=True) trigger_type = Column(String, default="manual") - events = relationship("JobEvent", back_populates="job", cascade="all, delete-orphan") + events = relationship( + "JobEvent", back_populates="job", cascade="all, delete-orphan" + ) class JobEvent(Base): __tablename__ = "job_events" id = Column(String, primary_key=True, default=lambda: str(uuid4())) - job_id = Column(String, ForeignKey("jobs.id", ondelete="CASCADE"), nullable=False, index=True) + job_id = Column( + String, ForeignKey("jobs.id", ondelete="CASCADE"), nullable=False, index=True + ) phase = Column(String, nullable=False) message = Column(Text, default="") - timestamp = Column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc)) + timestamp = Column(DateTime(timezone=True), default=lambda: datetime.now(UTC)) job = relationship("Job", back_populates="events") @@ -78,7 +85,9 @@ class JobEvent(Base): # Engine & session factory # --------------------------------------------------------------------------- -engine = create_async_engine(settings.database_url, echo=False, pool_size=5, max_overflow=10) +engine = create_async_engine( + settings.database_url, echo=False, pool_size=5, max_overflow=10 +) async_session = async_sessionmaker(engine, class_=AsyncSession, expire_on_commit=False) @@ -93,12 +102,15 @@ async def init_db() -> None: # Query helpers # --------------------------------------------------------------------------- + async def upsert_job(session: AsyncSession, **kwargs: Any) -> Job: """Insert or update a job record keyed by name (atomic).""" if "id" not in kwargs: kwargs["id"] = str(uuid4()) - update_cols = {k: v for k, v in kwargs.items() if k not in ("id", "name") and v is not None} + update_cols = { + k: v for k, v in kwargs.items() if k not in ("id", "name") and v is not None + } stmt = ( pg_insert(Job) @@ -165,7 +177,8 @@ async def list_jobs( async def list_jobs_by_schedule( - session: AsyncSession, schedule_name: str, + session: AsyncSession, + schedule_name: str, ) -> Sequence[Job]: """List all jobs triggered by a specific schedule.""" result = await session.execute( diff --git a/fournos-ui/app/forge_discovery.py b/fournos-ui/app/forge_discovery.py index 99715ff..2a7ee60 100644 --- a/fournos-ui/app/forge_discovery.py +++ b/fournos-ui/app/forge_discovery.py @@ -50,7 +50,11 @@ def _discover_from_repo(projects_dir: Path) -> dict[str, ProjectInfo]: skip = {"core", "__pycache__"} for proj_dir in sorted(projects_dir.iterdir()): - if not proj_dir.is_dir() or proj_dir.name.startswith(".") or proj_dir.name in skip: + if ( + not proj_dir.is_dir() + or proj_dir.name.startswith(".") + or proj_dir.name in skip + ): continue orchestration = proj_dir / "orchestration" @@ -92,7 +96,11 @@ def _discover_from_configmap() -> dict[str, ProjectInfo]: result: dict[str, ProjectInfo] = {} for idx, proj in enumerate(data["projects"]): if not isinstance(proj, dict): - logger.warning("Skipping malformed project entry at index %d: expected mapping, got %s", idx, type(proj).__name__) + logger.warning( + "Skipping malformed project entry at index %d: expected mapping, got %s", + idx, + type(proj).__name__, + ) continue name = proj.get("name", "") if not name: diff --git a/fournos-ui/app/k8s_client.py b/fournos-ui/app/k8s_client.py index abfce55..262dee8 100644 --- a/fournos-ui/app/k8s_client.py +++ b/fournos-ui/app/k8s_client.py @@ -5,10 +5,10 @@ import json import logging import threading -from datetime import datetime, timezone -from typing import Any, Generator +from collections.abc import Generator +from datetime import UTC, datetime +from typing import Any -import yaml from kubernetes import client, config, watch from kubernetes.client.rest import ApiException @@ -64,6 +64,7 @@ def is_connected() -> bool: # FournosJob operations # --------------------------------------------------------------------------- + def list_fournos_jobs(namespace: str | None = None) -> list[dict]: """List all FournosJob CRs in the given namespace.""" _ensure_loaded() @@ -119,9 +120,7 @@ def create_fournos_job(body: dict, namespace: str | None = None) -> dict: ) -def patch_fournos_job( - name: str, patch: dict, namespace: str | None = None -) -> dict: +def patch_fournos_job(name: str, patch: dict, namespace: str | None = None) -> dict: """Patch a FournosJob (e.g. set spec.shutdown).""" _ensure_loaded() if _custom_api is None: @@ -166,8 +165,7 @@ def watch_fournos_jobs( if timeout: kwargs["timeout_seconds"] = timeout try: - for event in w.stream(_custom_api.list_namespaced_custom_object, **kwargs): - yield event + yield from w.stream(_custom_api.list_namespaced_custom_object, **kwargs) except ApiException as exc: logger.warning("Watch stream ended: %s", exc.reason) @@ -176,6 +174,7 @@ def watch_fournos_jobs( # Tekton PipelineRun operations # --------------------------------------------------------------------------- + def get_pipelinerun(name: str, namespace: str | None = None) -> dict | None: """Get a Tekton PipelineRun by name.""" _ensure_loaded() @@ -317,14 +316,16 @@ def extract_pipeline_stages(pipelinerun: dict) -> list[dict]: display_name = task_name.replace("-", " ").title() - stages.append({ - "name": task_name, - "displayName": display_name, - "status": task_phase, - "startTime": start_time, - "completionTime": completion_time, - "finally": task_name in finally_task_names, - }) + stages.append( + { + "name": task_name, + "displayName": display_name, + "status": task_phase, + "startTime": start_time, + "completionTime": completion_time, + "finally": task_name in finally_task_names, + } + ) stages.sort(key=lambda s: (s["finally"], s.get("startTime") or "9999")) return stages @@ -334,6 +335,7 @@ def extract_pipeline_stages(pipelinerun: dict) -> list[dict]: # Pod operations # --------------------------------------------------------------------------- + def list_pods_for_job(job_name: str, namespace: str | None = None) -> list[dict]: """List pods associated with a FournosJob.""" _ensure_loaded() @@ -350,7 +352,7 @@ def list_pods_for_job(job_name: str, namespace: str | None = None) -> list[dict] created = pod.metadata.creation_timestamp age_minutes = 0 if created: - delta = datetime.now(timezone.utc) - created.replace(tzinfo=timezone.utc) + delta = datetime.now(UTC) - created.replace(tzinfo=UTC) age_minutes = int(delta.total_seconds() / 60) container_ready = False @@ -364,18 +366,22 @@ def list_pods_for_job(job_name: str, namespace: str | None = None) -> list[dict] if pod.metadata.name.startswith("affinity-assistant"): continue - pods.append({ - "name": pod.metadata.name, - "phase": pod.status.phase or "Unknown", - "container": ( - pod.spec.containers[0].name if pod.spec.containers else "unknown" - ), - "ready": container_ready, - "restarts": restarts, - "age_minutes": age_minutes, - "_created": created, - }) - pods.sort(key=lambda p: p["_created"] or datetime.min.replace(tzinfo=timezone.utc)) + pods.append( + { + "name": pod.metadata.name, + "phase": pod.status.phase or "Unknown", + "container": ( + pod.spec.containers[0].name + if pod.spec.containers + else "unknown" + ), + "ready": container_ready, + "restarts": restarts, + "age_minutes": age_minutes, + "_created": created, + } + ) + pods.sort(key=lambda p: p["_created"] or datetime.min.replace(tzinfo=UTC)) return pods except ApiException as exc: logger.error("Failed to list pods for %s: %s", job_name, exc.reason) @@ -402,7 +408,9 @@ def read_pod_log( kwargs["tail_lines"] = tail_lines try: if follow: - for line in _core_api.read_namespaced_pod_log(**kwargs, _preload_content=False).stream(): + for line in _core_api.read_namespaced_pod_log( + **kwargs, _preload_content=False + ).stream(): decoded = line.decode("utf-8", errors="replace").rstrip("\n") yield decoded else: @@ -601,7 +609,9 @@ def create_cronjob( if resolver_script: is_python = resolver_filename.lower().endswith(".py") default_resolver_img = "python:3.12-slim" if is_python else "alpine:latest" - script_key = resolver_filename or ("resolver.py" if is_python else "resolver.sh") + script_key = resolver_filename or ( + "resolver.py" if is_python else "resolver.sh" + ) configmap_name = f"{name}-resolver" _create_resolver_configmap(configmap_name, script_key, resolver_script, ns) @@ -619,7 +629,9 @@ def create_cronjob( ) volumes = [script_vol, shared_vol] script_mount = client.V1VolumeMount( - name="resolver-script", mount_path="/resolver", read_only=True, + name="resolver-script", + mount_path="/resolver", + read_only=True, ) shared_mount = client.V1VolumeMount(name="shared", mount_path="/shared") submit_container.volume_mounts = [shared_mount] @@ -714,13 +726,17 @@ def _create_resolver_configmap( except ApiException as exc: if exc.status == 409: _core_api.replace_namespaced_config_map( - name=cm_name, namespace=namespace, body=cm, + name=cm_name, + namespace=namespace, + body=cm, ) else: raise -def get_resolver_script(configmap_name: str, namespace: str | None = None) -> tuple[str, str]: +def get_resolver_script( + configmap_name: str, namespace: str | None = None +) -> tuple[str, str]: """Read the resolver script from its ConfigMap. Returns (filename, content).""" _ensure_loaded() if _core_api is None: @@ -755,13 +771,17 @@ def trigger_cronjob(name: str, namespace: str | None = None) -> str: ns = namespace or settings.fournos_namespace cj = _batch_api.read_namespaced_cron_job(name=name, namespace=ns) - ts = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + ts = datetime.now(UTC).strftime("%Y%m%d-%H%M%S") job_name = f"{name}-manual-{ts}"[:63] job_spec = cj.spec.job_template.spec trigger_env = client.V1EnvVar(name="FOURNOS_TRIGGER_TYPE", value="manual") - if job_spec.template and job_spec.template.spec and job_spec.template.spec.containers: + if ( + job_spec.template + and job_spec.template.spec + and job_spec.template.spec.containers + ): for c in job_spec.template.spec.containers: if c.env is None: c.env = [] @@ -839,7 +859,9 @@ def _cronjob_to_dict(cj: Any) -> dict: "resolver_image": resolver_image, "resolver_filename": resolver_filename, "has_resolver": bool(resolver_configmap), - "created_at": meta.creation_timestamp.isoformat() if meta.creation_timestamp else "", + "created_at": meta.creation_timestamp.isoformat() + if meta.creation_timestamp + else "", "last_schedule": ( cj.status.last_schedule_time.isoformat() if cj.status and cj.status.last_schedule_time diff --git a/fournos-ui/app/main.py b/fournos-ui/app/main.py index 72e51a0..632957f 100644 --- a/fournos-ui/app/main.py +++ b/fournos-ui/app/main.py @@ -6,7 +6,7 @@ import logging import re from contextlib import asynccontextmanager -from datetime import datetime, timezone +from datetime import UTC, datetime from pathlib import Path from typing import Any @@ -17,7 +17,7 @@ from app import db, k8s_client, watcher from app.config import settings -from app.forge_discovery import discover_projects, get_project_presets +from app.forge_discovery import discover_projects logger = logging.getLogger(__name__) @@ -25,6 +25,7 @@ # Lifespan # --------------------------------------------------------------------------- + @asynccontextmanager async def lifespan(app: FastAPI): logging.basicConfig(level=getattr(logging, settings.log_level)) @@ -32,6 +33,7 @@ async def lifespan(app: FastAPI): watcher.start_watcher() yield + app = FastAPI(title="Fournos Launcher Dashboard", lifespan=lifespan) BASE_DIR = Path(__file__).resolve().parent @@ -47,6 +49,7 @@ async def lifespan(app: FastAPI): # Template helpers # --------------------------------------------------------------------------- + def _format_age(timestamp_str: str) -> str: from dateutil.parser import parse @@ -54,7 +57,7 @@ def _format_age(timestamp_str: str) -> str: created = parse(timestamp_str) except Exception: return "?" - delta = datetime.now(timezone.utc) - created + delta = datetime.now(UTC) - created total_seconds = int(delta.total_seconds()) if total_seconds < 0: return "0s" @@ -97,7 +100,11 @@ def _extract_forge_info(job: dict) -> dict: pr_title = env.get("PULL_TITLE", "") repo_owner = env.get("REPO_OWNER", "") repo_name = env.get("REPO_NAME", "") - pr_url = f"https://github.com/{repo_owner}/{repo_name}/pull/{pr_number}" if pr_number else "" + pr_url = ( + f"https://github.com/{repo_owner}/{repo_name}/pull/{pr_number}" + if pr_number + else "" + ) return { "project": forge.get("project", ""), "args": forge.get("args", []), @@ -128,7 +135,7 @@ def _parse_task_progress(message: str) -> dict | None: def _build_timeline(stages: list[dict]) -> list[dict]: from dateutil.parser import parse - now = datetime.now(timezone.utc) + now = datetime.now(UTC) n = len(stages) or 1 equal_pct = 100.0 / n @@ -160,13 +167,15 @@ def _build_timeline(stages: list[dict]) -> list[dict]: "Skipped": "ptl-skip", }.get(s["status"], "ptl-wait") - result.append({ - **s, - "width_pct": equal_pct, - "min_width": 8, - "duration_label": dur_label if s["startTime"] else "", - "status_class": status_class, - }) + result.append( + { + **s, + "width_pct": equal_pct, + "min_width": 8, + "duration_label": dur_label if s["startTime"] else "", + "status_class": status_class, + } + ) return result @@ -183,7 +192,7 @@ def _extract_mlflow_url(status: dict) -> str: return mlflow.get("run_url", "") if mlflow else "" -_CACHE_BUST = str(int(datetime.now(timezone.utc).timestamp())) +_CACHE_BUST = str(int(datetime.now(UTC).timestamp())) _jinja_env.globals.update( format_age=_format_age, @@ -226,7 +235,7 @@ def _get_live_jobs_sync() -> list[dict]: from dateutil.parser import parse jobs = k8s_client.list_fournos_jobs() - now = datetime.now(timezone.utc) + now = datetime.now(UTC) visible: list[dict] = [] for j in jobs: phase = j.get("status", {}).get("phase", "") @@ -304,6 +313,7 @@ async def _get_pipeline_stages(job: dict) -> list[dict]: # Routes: Jobs # --------------------------------------------------------------------------- + @app.get("/", response_class=HTMLResponse) async def jobs_list( request: Request, @@ -330,7 +340,7 @@ async def jobs_list( total = len(jobs) all_clusters = _collect_clusters(jobs) offset = (page - 1) * per_page - jobs = jobs[offset:offset + per_page] + jobs = jobs[offset : offset + per_page] history_jobs = [] total_history = 0 else: @@ -386,7 +396,9 @@ async def jobs_table_partial( if owner: jobs = [j for j in jobs if j.get("spec", {}).get("owner") == owner] current_steps = await _compute_current_steps(jobs) - return _render("components/jobs_table_body.html", jobs=jobs, current_steps=current_steps) + return _render( + "components/jobs_table_body.html", jobs=jobs, current_steps=current_steps + ) @app.get("/jobs/{job_name}", response_class=HTMLResponse) @@ -471,7 +483,11 @@ async def rerun_job(job_name: str): try: created = await asyncio.to_thread(k8s_client.create_fournos_job, body) created_name = created.get("metadata", {}).get("name", new_name) - return {"status": "ok", "job_name": created_name, "redirect": f"/jobs/{created_name}"} + return { + "status": "ok", + "job_name": created_name, + "redirect": f"/jobs/{created_name}", + } except Exception as exc: raise HTTPException(status_code=500, detail=str(exc)) @@ -491,9 +507,8 @@ async def _get_job_for_rerun(job_name: str) -> dict | None: @app.delete("/api/history/{job_name}") async def delete_history_job(job_name: str): """Delete a job from the history database.""" - async with db.async_session() as session: - async with session.begin(): - deleted = await db.delete_job_by_name(session, job_name) + async with db.async_session() as session, session.begin(): + deleted = await db.delete_job_by_name(session, job_name) if not deleted: raise HTTPException(status_code=404, detail="Job not found in history") return {"status": "ok"} @@ -524,7 +539,7 @@ def _reader(): finally: loop.call_soon_threadsafe(queue.put_nowait, None) - task = asyncio.get_event_loop().run_in_executor(None, _reader) + asyncio.get_event_loop().run_in_executor(None, _reader) try: while True: line = await queue.get() @@ -537,12 +552,11 @@ def _reader(): return StreamingResponse(generate(), media_type="text/event-stream") - - # --------------------------------------------------------------------------- # Routes: Submit Job # --------------------------------------------------------------------------- + @app.get("/submit", response_class=HTMLResponse) async def submit_form(request: Request): projects = discover_projects() @@ -556,6 +570,7 @@ async def submit_form(request: Request): @app.get("/api/project-info/{project_name}") async def project_info_api(project_name: str): from app.forge_discovery import get_project + proj = get_project(project_name) if proj is None: return {"presets": [], "cluster": ""} @@ -564,8 +579,8 @@ async def project_info_api(project_name: str): def _fetch_github_open_prs() -> list[dict]: """Blocking call to the GitHub API -- run via asyncio.to_thread.""" - import urllib.request import json as _json + import urllib.request url = f"https://api.github.com/repos/{settings.forge_github_repo}/pulls?state=open&per_page=100" req = urllib.request.Request(url, headers={"Accept": "application/vnd.github+json"}) @@ -670,22 +685,25 @@ async def submit_job( created_name = created.get("metadata", {}).get("name", job_name) try: - async with db.async_session() as session: - async with session.begin(): - await db.upsert_job( - session, - name=created_name, - project=project, - preset=preset, - cluster=cluster, - pipeline=pipeline, - owner=owner or "fournos-dashboard", - status="Pending", - config_overrides=config_overrides, - fjob_spec=body.get("spec", {}), - ) + async with db.async_session() as session, session.begin(): + await db.upsert_job( + session, + name=created_name, + project=project, + preset=preset, + cluster=cluster, + pipeline=pipeline, + owner=owner or "fournos-dashboard", + status="Pending", + config_overrides=config_overrides, + fjob_spec=body.get("spec", {}), + ) except Exception as exc: - logger.error("DB upsert failed for job %s (job was created in K8s): %s", created_name, exc) + logger.error( + "DB upsert failed for job %s (job was created in K8s): %s", + created_name, + exc, + ) return RedirectResponse(url=f"/jobs/{created_name}", status_code=303) @@ -694,6 +712,7 @@ async def submit_job( # Routes: Schedules # --------------------------------------------------------------------------- + @app.get("/schedules", response_class=HTMLResponse) async def schedules_list(request: Request): cronjobs = await asyncio.to_thread(k8s_client.list_managed_cronjobs) @@ -713,15 +732,17 @@ async def schedule_runs(request: Request, name: str): jobs = await db.list_jobs_by_schedule(session, name) runs = [] for j in jobs: - runs.append({ - "name": j.name, - "status": j.status, - "preset": j.preset, - "trigger_type": j.trigger_type or "scheduled", - "duration_seconds": j.duration_seconds, - "mlflow_url": j.mlflow_url, - "created_at": j.created_at.isoformat() if j.created_at else "", - }) + runs.append( + { + "name": j.name, + "status": j.status, + "preset": j.preset, + "trigger_type": j.trigger_type or "scheduled", + "duration_seconds": j.duration_seconds, + "mlflow_url": j.mlflow_url, + "created_at": j.created_at.isoformat() if j.created_at else "", + } + ) return _render("schedule_runs.html", schedule_name=name, runs=runs) @@ -753,14 +774,20 @@ async def create_schedule( preset=preset, image=image_source, owner=owner, - resolver_script=resolver_script.strip().replace("\r\n", "\n").replace("\r", "\n"), + resolver_script=resolver_script.strip() + .replace("\r\n", "\n") + .replace("\r", "\n"), resolver_image=resolver_image.strip(), resolver_filename=resolver_filename.strip(), ) try: await asyncio.to_thread(k8s_client.delete_cronjob, edit_target) except Exception as del_exc: - logger.warning("Failed to delete old schedule %s after replacement: %s", edit_target, del_exc) + logger.warning( + "Failed to delete old schedule %s after replacement: %s", + edit_target, + del_exc, + ) else: if edit_target: await asyncio.to_thread(k8s_client.delete_cronjob, edit_target) @@ -774,7 +801,9 @@ async def create_schedule( preset=preset, image=image_source, owner=owner, - resolver_script=resolver_script.strip().replace("\r\n", "\n").replace("\r", "\n"), + resolver_script=resolver_script.strip() + .replace("\r\n", "\n") + .replace("\r", "\n"), resolver_image=resolver_image.strip(), resolver_filename=resolver_filename.strip(), ) @@ -839,6 +868,7 @@ async def delete_schedule(name: str): # Conversion helpers # --------------------------------------------------------------------------- + def _collect_clusters(live_jobs: list[dict]) -> list[str]: """Collect unique cluster names from live jobs.""" clusters = set() @@ -878,7 +908,11 @@ def _db_job_to_fjob_dict(job: db.Job) -> dict: forge = spec.get("executionEngine", {}).get("forge", {}) if not forge: - forge = {"project": job.project, "args": job.preset.split() if job.preset else [], "configOverrides": job.config_overrides or {}} + forge = { + "project": job.project, + "args": job.preset.split() if job.preset else [], + "configOverrides": job.config_overrides or {}, + } spec.setdefault("executionEngine", {})["forge"] = forge spec.setdefault("cluster", job.cluster) @@ -921,10 +955,8 @@ def _get_version_config_key(project: str) -> str: def sanitize_job_name(prefix: str) -> str: """Generate a K8s-safe job name with timestamp.""" - ts = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + ts = datetime.now(UTC).strftime("%Y%m%d-%H%M%S") name = f"{prefix}-{ts}".lower() name = re.sub(r"[^a-z0-9-]", "-", name) name = re.sub(r"-+", "-", name).strip("-") return name[:63] - - diff --git a/fournos-ui/app/watcher.py b/fournos-ui/app/watcher.py index fde6429..7721d96 100644 --- a/fournos-ui/app/watcher.py +++ b/fournos-ui/app/watcher.py @@ -6,7 +6,6 @@ import logging import threading import time -from datetime import datetime, timezone from dateutil.parser import parse as parse_dt from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine @@ -27,10 +26,15 @@ def _init_watcher_db(loop: asyncio.AbstractEventLoop) -> None: """Create a separate DB engine for the watcher's own event loop.""" global _watcher_engine, _watcher_session _watcher_engine = create_async_engine( - settings.database_url, echo=False, pool_size=3, max_overflow=5, + settings.database_url, + echo=False, + pool_size=3, + max_overflow=5, ) _watcher_session = async_sessionmaker( - _watcher_engine, class_=AsyncSession, expire_on_commit=False, + _watcher_engine, + class_=AsyncSession, + expire_on_commit=False, ) @@ -66,7 +70,10 @@ def _extract_forge_fields(job: dict) -> dict: duration_seconds = None conditions = status.get("conditions", []) for cond in conditions: - if cond.get("type") == "PipelineRunReady" and cond.get("status") in ("True", "False"): + if cond.get("type") == "PipelineRunReady" and cond.get("status") in ( + "True", + "False", + ): try: completed_at = parse_dt(cond["lastTransitionTime"]) except Exception: @@ -111,21 +118,20 @@ async def _archive_job(job: dict) -> None: logger.warning("Skipping FournosJob with missing name") return - async with _watcher_session() as session: - async with session.begin(): - existing = await db.get_job_by_name(session, job_name) - previous_phase = existing.status if existing else None - previous_message = existing.message if existing else None + async with _watcher_session() as session, session.begin(): + existing = await db.get_job_by_name(session, job_name) + previous_phase = existing.status if existing else None + previous_message = existing.message if existing else None - db_job = await db.upsert_job(session, **fields) + db_job = await db.upsert_job(session, **fields) - if fields["status"] != previous_phase or fields["message"] != previous_message: - await db.add_job_event( - session, - job_id=db_job.id, - phase=fields["status"], - message=fields["message"], - ) + if fields["status"] != previous_phase or fields["message"] != previous_message: + await db.add_job_event( + session, + job_id=db_job.id, + phase=fields["status"], + message=fields["message"], + ) logger.info("Archived FournosJob %s (phase=%s)", job_name, fields["status"]) @@ -159,7 +165,12 @@ async def _full_sync() -> None: logger.warning("Full sync: failed to archive %s: %s", name, exc) errors += 1 - logger.info("Full sync complete: %d synced, %d errors (out of %d)", synced, errors, len(all_jobs)) + logger.info( + "Full sync complete: %d synced, %d errors (out of %d)", + synced, + errors, + len(all_jobs), + ) def _run_watch_loop(loop: asyncio.AbstractEventLoop) -> None: @@ -184,7 +195,9 @@ def _run_watch_loop(loop: asyncio.AbstractEventLoop) -> None: while True: try: - logger.info("Starting FournosJob watch (rv=%s)", resource_version or "latest") + logger.info( + "Starting FournosJob watch (rv=%s)", resource_version or "latest" + ) for event in k8s_client.watch_fournos_jobs( resource_version=resource_version, timeout=300, @@ -202,7 +215,9 @@ def _run_watch_loop(loop: asyncio.AbstractEventLoop) -> None: name = obj.get("metadata", {}).get("name", "?") logger.error( "Failed to archive event for %s (type=%s): %s", - name, event_type, exc, + name, + event_type, + exc, ) elif event_type == "DELETED": name = obj.get("metadata", {}).get("name", "") From 21acf82dc2f4f078be88f83b81b32d96f3997f4b Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:13:35 +0200 Subject: [PATCH 2/6] founos: lint --- fournos-ui/app/config.py | 16 ++++++++-------- fournos/__init__.py | 2 +- fournos/__main__.py | 2 +- fournos/handlers/__init__.py | 6 +++--- fournos/handlers/lifecycle.py | 2 +- fournos/handlers/status.py | 2 +- fournos/operator.py | 2 +- 7 files changed, 16 insertions(+), 16 deletions(-) diff --git a/fournos-ui/app/config.py b/fournos-ui/app/config.py index aeef2e2..e463895 100644 --- a/fournos-ui/app/config.py +++ b/fournos-ui/app/config.py @@ -8,9 +8,7 @@ @dataclass(frozen=True) class Settings: - database_url: str = field( - default_factory=lambda: os.environ["DATABASE_URL"] - ) + database_url: str = field(default_factory=lambda: os.environ["DATABASE_URL"]) fournos_namespace: str = field( default_factory=lambda: os.environ.get("FOURNOS_NAMESPACE", "fournos-jobs") @@ -25,7 +23,9 @@ class Settings: ) projects_config_path: str = field( - default_factory=lambda: os.environ.get("PROJECTS_CONFIG_PATH", "/etc/fournos-dashboard/projects.yaml") + default_factory=lambda: os.environ.get( + "PROJECTS_CONFIG_PATH", "/etc/fournos-dashboard/projects.yaml" + ) ) fournos_api_group: str = "fournos.dev" @@ -36,12 +36,12 @@ class Settings: tekton_api_version: str = "v1" tekton_pipelinerun_plural: str = "pipelineruns" - log_level: str = field( - default_factory=lambda: os.environ.get("LOG_LEVEL", "INFO") - ) + log_level: str = field(default_factory=lambda: os.environ.get("LOG_LEVEL", "INFO")) forge_github_repo: str = field( - default_factory=lambda: os.environ.get("FORGE_GITHUB_REPO", "openshift-psap/forge") + default_factory=lambda: os.environ.get( + "FORGE_GITHUB_REPO", "openshift-psap/forge" + ) ) k8s_request_timeout_seconds: int = field( diff --git a/fournos/__init__.py b/fournos/__init__.py index 0e1d312..068c4c1 100644 --- a/fournos/__init__.py +++ b/fournos/__init__.py @@ -1,4 +1,4 @@ -from importlib.metadata import version, PackageNotFoundError +from importlib.metadata import PackageNotFoundError, version try: __version__ = version("fournos") diff --git a/fournos/__main__.py b/fournos/__main__.py index 961e248..87ade23 100644 --- a/fournos/__main__.py +++ b/fournos/__main__.py @@ -19,6 +19,6 @@ *sys.argv[1:], ] -from kopf.cli import main # noqa: E402 +from kopf.cli import main main() diff --git a/fournos/handlers/__init__.py b/fournos/handlers/__init__.py index f4711af..e6be78d 100644 --- a/fournos/handlers/__init__.py +++ b/fournos/handlers/__init__.py @@ -2,19 +2,19 @@ from .execution import ( handle_shutdown, - reconcile_stopping, reconcile_admitted, reconcile_running, + reconcile_stopping, ) from .lifecycle import on_create, reconcile_pending from .resolving import reconcile_resolving __all__ = [ + "handle_shutdown", "on_create", + "reconcile_admitted", "reconcile_pending", "reconcile_resolving", - "reconcile_admitted", "reconcile_running", - "handle_shutdown", "reconcile_stopping", ] diff --git a/fournos/handlers/lifecycle.py b/fournos/handlers/lifecycle.py index c10eb8b..12a2444 100644 --- a/fournos/handlers/lifecycle.py +++ b/fournos/handlers/lifecycle.py @@ -22,9 +22,9 @@ from fournos.state import ctx from .status import ( + COND_WORKLOAD_ADMITTED, CRD_GROUP, CRD_VERSION, - COND_WORKLOAD_ADMITTED, owner_ref, set_condition, ) diff --git a/fournos/handlers/status.py b/fournos/handlers/status.py index 6dabacb..e102436 100644 --- a/fournos/handlers/status.py +++ b/fournos/handlers/status.py @@ -25,7 +25,7 @@ def owner_ref(body: dict) -> dict: def utcnow() -> str: - return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + return datetime.datetime.now(datetime.UTC).strftime("%Y-%m-%dT%H:%M:%SZ") def set_condition( diff --git a/fournos/operator.py b/fournos/operator.py index 98ecf99..fb242a4 100644 --- a/fournos/operator.py +++ b/fournos/operator.py @@ -13,12 +13,12 @@ import kopf from kubernetes import client, config +from fournos import __version__, handlers from fournos.core.clusters import ClusterRegistry from fournos.core.constants import LABEL_JOB_NAME, Phase from fournos.core.kueue import KueueClient from fournos.core.resolve import ResolveClient from fournos.core.tekton import TektonClient -from fournos import __version__, handlers from fournos.settings import settings from fournos.state import ctx From f1cba3ad6aed220999ce2205f719e095633f9332 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:39:25 +0200 Subject: [PATCH 3/6] hacks: sync_vault_secrets: lint --- hacks/sync_vault_secrets.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/hacks/sync_vault_secrets.py b/hacks/sync_vault_secrets.py index d43a712..99bf354 100755 --- a/hacks/sync_vault_secrets.py +++ b/hacks/sync_vault_secrets.py @@ -132,7 +132,8 @@ def vault_read( def _k8s_core_api(): - from kubernetes import client, config as k8s_config + from kubernetes import client + from kubernetes import config as k8s_config try: k8s_config.load_kube_config() @@ -326,8 +327,8 @@ def sync( if dry_run: print(f"[dry-run] Would create/update Secret {namespace}/{secret_name}") - for key in safe_data: - print(f" {key}: <{len(str(safe_data[key]))} chars>") + for key, value in safe_data.items(): + print(f" {key}: <{len(str(value))} chars>") processed_secrets.add(secret_name) continue From 545355ff41cbd5f717ad92a1fc1ef18a29441515 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:39:43 +0200 Subject: [PATCH 4/6] tests: forge: lint --- tests/forge/deploy/orchestration/ci.py | 18 +++++++++--------- tests/forge/deploy/orchestration/cli.py | 10 +++++----- tests/forge/deploy/orchestration/deploy.py | 12 ++++++------ tests/forge/deploy/orchestration/utils.py | 4 ++-- 4 files changed, 22 insertions(+), 22 deletions(-) diff --git a/tests/forge/deploy/orchestration/ci.py b/tests/forge/deploy/orchestration/ci.py index 6baa2eb..02dc205 100755 --- a/tests/forge/deploy/orchestration/ci.py +++ b/tests/forge/deploy/orchestration/ci.py @@ -3,17 +3,17 @@ Fournos-Deploy Project CI Operations """ -from projects.core.library import ci as ci_lib, config, env -from projects.core.ci_entrypoint.prepare_ci import CI_METADATA_DIRNAME -from projects.fournos_launcher.orchestration import utils - -import deploy as fournos_deploy - -import click -import types +import json import logging import os -import json +import types + +import click +import deploy as fournos_deploy +from projects.core.ci_entrypoint.prepare_ci import CI_METADATA_DIRNAME +from projects.core.library import ci as ci_lib +from projects.core.library import config, env +from projects.fournos_launcher.orchestration import utils logger = logging.getLogger(__name__) diff --git a/tests/forge/deploy/orchestration/cli.py b/tests/forge/deploy/orchestration/cli.py index 0fabb52..5a9b765 100755 --- a/tests/forge/deploy/orchestration/cli.py +++ b/tests/forge/deploy/orchestration/cli.py @@ -5,14 +5,14 @@ Interactive CLI for FOURNOS deployment with configuration overrides. """ -import deploy as fournos_deploy -from projects.core.library.cli import safe_cli_command -from projects.core.library import config, run - +import logging import sys import types + import click -import logging +import deploy as fournos_deploy +from projects.core.library import config, run +from projects.core.library.cli import safe_cli_command logger = logging.getLogger(__name__) diff --git a/tests/forge/deploy/orchestration/deploy.py b/tests/forge/deploy/orchestration/deploy.py index 9ec8a5b..77861c6 100644 --- a/tests/forge/deploy/orchestration/deploy.py +++ b/tests/forge/deploy/orchestration/deploy.py @@ -1,12 +1,12 @@ -from projects.core.library import env, config, run, vault -from projects.cluster.toolbox.build_image.main import run as build_image_toolbox - -import pathlib import logging import os -import yaml +import pathlib from pathlib import Path +import yaml +from projects.cluster.toolbox.build_image.main import run as build_image_toolbox +from projects.core.library import config, env, run, vault + logger = logging.getLogger(__name__) @@ -46,7 +46,7 @@ def _apply_manifest_replacements(manifest_file): # Prepare replacements with config value resolution resolved_replacements = {} - for key in replacements_config.keys(): + for key in replacements_config: resolved_value = config.project.get_config( f"fournos_deploy.manifests.replace.{key}", print=False ) diff --git a/tests/forge/deploy/orchestration/utils.py b/tests/forge/deploy/orchestration/utils.py index 2d1ceff..435afa9 100644 --- a/tests/forge/deploy/orchestration/utils.py +++ b/tests/forge/deploy/orchestration/utils.py @@ -1,10 +1,10 @@ +import logging +import os import shutil import tarfile import tempfile import urllib.request from pathlib import Path -import os -import logging logger = logging.getLogger(__name__) From 11b78802d75756a95055e75173707604cfbb18d8 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:45:06 +0200 Subject: [PATCH 5/6] tests: conftest: lint --- tests/conftest.py | 4 ++++ tests/test_exclusive.py | 1 + 2 files changed, 5 insertions(+) diff --git a/tests/conftest.py b/tests/conftest.py index 63c2bf9..2f27901 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -181,6 +181,7 @@ def workload_exists(name: str) -> bool: result = subprocess.run( ["kubectl", "get", "workload", name, "-n", NAMESPACE], capture_output=True, + check=False, ) return result.returncode == 0 @@ -190,6 +191,7 @@ def pipelinerun_exists(name: str) -> bool: result = subprocess.run( ["kubectl", "get", "pipelinerun", name, "-n", NAMESPACE], capture_output=True, + check=False, ) return result.returncode == 0 @@ -299,6 +301,7 @@ def resolve_job_exists(name: str) -> bool: result = subprocess.run( ["kubectl", "get", "job", f"{name}-resolve", "-n", NAMESPACE], capture_output=True, + check=False, ) return result.returncode == 0 @@ -329,6 +332,7 @@ def poll_resolve_job_complete( ], capture_output=True, text=True, + check=False, ) conditions = result.stdout.strip() if "Complete" in conditions or "Failed" in conditions: diff --git a/tests/test_exclusive.py b/tests/test_exclusive.py index 13bd3ae..45f85a2 100644 --- a/tests/test_exclusive.py +++ b/tests/test_exclusive.py @@ -42,6 +42,7 @@ def _slow_mock_pipeline(): ], capture_output=True, text=True, + check=False, ) original_sleep = result.stdout.strip() or "3" From b678fc97d079665e646c95f25b9ab87271983773 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Thu, 30 Jul 2026 20:50:47 +0200 Subject: [PATCH 6/6] pyproject: ruff: ignore BLE001 and S110 --- pyproject.toml | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/pyproject.toml b/pyproject.toml index d0a01ba..be785f0 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -26,3 +26,9 @@ testpaths = ["tests"] markers = [ "slow: tests that take longer than 30s (e.g. full lifecycle)", ] + +[tool.ruff.lint] +ignore = [ + "BLE001", # Do not catch blind exception: `Exception` + "S110", # `try`-`except`-`pass` detected, consider logging the exception +]