Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .claude-plugin/marketplace.json
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
"name": "edge-ocp-ci",
"source": "./plugins/edge-ocp-ci",
"description": "Edge OCP Payload Monitor — monitor OpenShift nightly payloads for edge topology (SNO/TNA/TNF) failures with AI-enriched analysis",
"version": "1.2.0"
"version": "1.2.1"
},
{
"name": "edge-scrum",
Expand Down
8 changes: 4 additions & 4 deletions payload-monitor/payload_monitor/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
from .collectors import component_readiness, prow, sippy, timing
from .collectors.release_controller import (
collect as collect_payloads,
discover_streams,
configured_stream_names,
version_from_stream,
)
from .config import Config
Expand Down Expand Up @@ -154,9 +154,9 @@ def main(

logger.info("Starting Edge OCP Payload Monitor")

# Step 1: Discover versions and resolve stream names
logger.info("Step 1: Discovering active versions...")
stream_names = discover_streams(config)
# Step 1: Resolve configured release streams
logger.info("Step 1: Resolving configured release streams...")
stream_names = configured_stream_names(config)
active_versions = [version_from_stream(s) for s in stream_names]
logger.info(f" Versions: {active_versions}")

Expand Down
17 changes: 14 additions & 3 deletions payload-monitor/payload_monitor/analyzer.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
)
from .collectors import jira as jira_collector
from .collectors.jira import has_auth as jira_has_auth
from .collectors.sippy import job_analysis_url

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -80,28 +81,38 @@ def _find_escalation_risks(

for job_name, topology in informing_jobs.items():
consecutive = 0
latest_prow_url = ""
Comment thread
vimauro marked this conversation as resolved.
latest_failure_seen = False
streak_runs: list[dict] = []
for payload in reversed_payloads:
job_in_payload = None
for job in payload.jobs:
if job.name == job_name:
job_in_payload = job
break
if job_in_payload is None:
# Job absent from this payload breaks the streak
break
if job_in_payload.result == JobResult.FAILURE:
consecutive += 1
streak_runs.append({
"payload_tag": payload.tag,
"prow_url": job_in_payload.prow_url,
})
if not latest_failure_seen:
latest_prow_url = job_in_payload.prow_url
latest_failure_seen = True
else:
break

if consecutive >= config.escalation_threshold:
sippy_url = f"https://sippy.dptools.openshift.org/sippy-ng/jobs/{job_name}"
risks.append(EscalationRisk(
job_name=job_name,
topology=topology,
version=stream.version,
consecutive_failures=consecutive,
sippy_url=sippy_url,
prow_url=latest_prow_url,
triage_url=job_analysis_url(stream.version, job_name),
failing_runs=streak_runs,
))
Comment thread
vimauro marked this conversation as resolved.

return risks
Expand Down
10 changes: 7 additions & 3 deletions payload-monitor/payload_monitor/collectors/release_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,8 +90,12 @@ def _parse_jobs(



def discover_streams(config: Config) -> list[str]:
"""Return list of nightly stream names to monitor."""
def configured_stream_names(config: Config) -> list[str]:
"""Map ``config.versions`` to release-controller stream names.

This is a static mapping, not a live discovery against the
release-controller index.
"""
return [_stream_name(v) for v in config.versions]


Expand Down Expand Up @@ -204,7 +208,7 @@ def _collect_stream(stream: str, config: Config) -> StreamReport:
def collect(config: Config, streams: Optional[list[str]] = None) -> list[StreamReport]:
"""Collect payload data for all configured streams in parallel."""
if streams is None:
streams = discover_streams(config)
streams = configured_stream_names(config)

with ThreadPoolExecutor(max_workers=min(len(streams) or 1, 10)) as pool:
futures = {
Expand Down
17 changes: 14 additions & 3 deletions payload-monitor/payload_monitor/collectors/sippy.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,10 @@

from __future__ import annotations

import json
import logging
from concurrent.futures import ThreadPoolExecutor, as_completed
from urllib.parse import quote

import requests

Expand Down Expand Up @@ -48,8 +50,17 @@ def fetch_edge_jobs(
return edge_jobs


def job_analysis_url(version: str, name: str) -> str:
"""Build the Sippy UI URL for a job's analysis page, filtered by job name."""
filters = json.dumps({
"items": [{"columnField": "name", "operatorValue": "equals", "value": name}]
})
return f"{BASE_URL}/sippy-ng/jobs/{version}/analysis?filters={quote(filters)}"


def identify_regressions(
edge_jobs: list[dict],
version: str = "",
min_runs: int = 3,
Comment thread
vimauro marked this conversation as resolved.
) -> list[Regression]:
"""Identify jobs that are getting worse over time.
Expand All @@ -69,9 +80,9 @@ def identify_regressions(
name = job.get("name", "")
topology = job.get("_topology", "")
jira_component = job.get("jira_component", "")
triage_url = f"{BASE_URL}/sippy-ng/jobs/{name}"
triage_url = job_analysis_url(version, name)

# Skip jobs with too few runs insufficient data to confirm regression
# Skip jobs with too few runs - insufficient data to confirm regression
if current_runs < min_runs:
continue

Expand Down Expand Up @@ -99,7 +110,7 @@ def _collect_version(version: str, config: Config) -> tuple[str, list[Regression
session = create_session()
try:
edge_jobs = fetch_edge_jobs(version, config, session=session)
regressions = identify_regressions(edge_jobs)
regressions = identify_regressions(edge_jobs, version=version)
if regressions:
logger.info(
f" {version}: {len(regressions)} edge job regressions detected"
Expand Down
27 changes: 6 additions & 21 deletions payload-monitor/payload_monitor/collectors/timing.py
Original file line number Diff line number Diff line change
Expand Up @@ -324,29 +324,14 @@ def compute_stats(runs: list[TimingRun]) -> dict:
# ---------------------------------------------------------------------------

def fetch_edge_jobs(release: str, config: Config) -> list[dict]:
"""Fetch all jobs from Sippy for a release, filter for edge topologies."""
logger.info(f"Fetching edge topology jobs for release {release}")
try:
resp = _session.get(JOBS_URL, params={"release": release}, timeout=30)
resp.raise_for_status()
all_jobs = resp.json()
except requests_lib.RequestException as e:
logger.error(f"Failed to fetch Sippy jobs for {release}: {e}")
return []
"""Fetch all jobs from Sippy for a release, filter for edge topologies.

if not isinstance(all_jobs, list):
return []

edge_jobs = []
for job in all_jobs:
name = job.get("name", "")
topology = config.classify_topology(name)
if topology in ("SNO", "TNA", "TNF"):
job["_topology"] = topology
edge_jobs.append(job)
Delegates to collectors.sippy.fetch_edge_jobs so topology filtering stays
config-driven and isn't duplicated across collectors.
"""
from . import sippy as _sippy

logger.info(f" Found {len(edge_jobs)} SNO/TNA/TNF jobs for {release}")
return edge_jobs
return _sippy.fetch_edge_jobs(release, config, session=_session)


def fetch_job_runs(job_name: str, release: str) -> list[dict]:
Expand Down
15 changes: 3 additions & 12 deletions payload-monitor/payload_monitor/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,9 @@ class EscalationRisk:
topology: str
version: str
consecutive_failures: int
sippy_url: str = ""
prow_url: str = ""
triage_url: str = ""
failing_runs: list[dict] = field(default_factory=list)


@dataclass
Expand Down Expand Up @@ -247,17 +249,6 @@ class TimingRun:
def is_success(self) -> bool:
return self.result == "S"

@property
def duration_minutes(self) -> float:
return self.duration_seconds / 60.0

@property
def install_duration_seconds(self) -> float:
"""Return install step duration if available, else 0."""
for key in ("install", "setup"):
if key in self.step_durations:
return self.step_durations[key]
return 0.0


@dataclass
Expand Down
Loading