This is an automated email from the ASF dual-hosted git repository.

potiuk pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/v3-3-test by this push:
     new a2d3a606d8c [v3-3-test] Discount CI image-build time from 
duration-trend alerts (#69789) (#70102)
a2d3a606d8c is described below

commit a2d3a606d8cd1c6796898aed8d5241d9f59292de
Author: github-actions[bot] 
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Mon Jul 20 00:33:08 2026 +0200

    [v3-3-test] Discount CI image-build time from duration-trend alerts 
(#69789) (#70102)
    
    * Discount CI image-build time from duration-trend alerts
    
    The CI duration monitor folded the "Prepare breeze & CI image" step into
    each job's total. That step occasionally balloons on a one-off cache miss
    (a full image rebuild instead of a cached pull), which repeatedly flagged
    unrelated jobs as regressed — e.g. a Non-DB core job reported +190% driven
    almost entirely by an 18m image build, while its actual test time was flat.
    
    Image-build time is now excluded from the run and per-job durations used
    for the trend, and the image build is watched on its own: it is only
    reported when it has stayed slow for more than two days, so a transient
    rebuild spike no longer produces a false alert.
    
    * Apply naming nits: prefix duration helpers with calculate_
    
    Rename work_duration -> calculate_work_duration,
    run_image_build_seconds -> calculate_image_build_seconds,
    adjusted_run_duration -> calculate_adjusted_run_duration for a
    consistent calculate_ prefix on the duration-computing helpers.
    
    
    
    ---------
    (cherry picked from commit 582076cc27396a01948b4c72fca006643c0846ea)
    
    Co-authored-by: Jarek Potiuk <[email protected]>
    Co-authored-by: Claude Opus 4.8 <[email protected]>
---
 scripts/ci/analyze_ci_job_durations.py            | 276 +++++++++++++++++-----
 scripts/tests/ci/test_analyze_ci_job_durations.py | 233 +++++++++++++++---
 2 files changed, 416 insertions(+), 93 deletions(-)

diff --git a/scripts/ci/analyze_ci_job_durations.py 
b/scripts/ci/analyze_ci_job_durations.py
index a03ddf7a57e..56f3c463333 100644
--- a/scripts/ci/analyze_ci_job_durations.py
+++ b/scripts/ci/analyze_ci_job_durations.py
@@ -35,6 +35,12 @@ both a relative margin (``REL_THRESHOLD``) and an absolute 
floor
 (``MIN_ABS_INCREASE_MINUTES`` / ``JOB_MIN_ABS_INCREASE_MINUTES``) so short jobs
 with noisy timings do not trigger spurious alerts.
 
+The image-build step ("Prepare breeze & CI image") occasionally balloons on a
+one-off cache miss, so its time is *excluded* from the run and per-job 
durations
+used for the trend above. The image build is instead watched on its own and 
only
+reported when it has stayed slow for longer than 
``IMAGE_BUILD_PERSISTENCE_DAYS``
+(so a single slow night never alerts).
+
 Environment variables (required):
   GITHUB_REPOSITORY  - Owner/repo (e.g. apache/airflow)
   GITHUB_TOKEN       - GitHub token for API access (used by ``gh``)
@@ -50,6 +56,8 @@ Environment variables (optional):
   REL_THRESHOLD             - Relative increase over baseline to flag, e.g. 
0.25 = 25% (default: 0.25)
   MIN_ABS_INCREASE_MINUTES  - Absolute floor for the overall-run alert 
(default: 5)
   JOB_MIN_ABS_INCREASE_MINUTES - Absolute floor for per-job alerts (default: 3)
+  IMAGE_BUILD_PERSISTENCE_DAYS - Only report a slow image build once it has 
stayed
+                              elevated for at least this many days (default: 2)
   ANALYZE_JOBS              - Whether to fetch per-job durations 
("true"/"false", default: true)
   ONLY_SUCCESSFUL           - Only consider runs that concluded "success" 
(default: true)
   SLACK_CHANNEL             - Slack channel for the message payload (default: 
internal-airflow-ci-cd)
@@ -185,18 +193,6 @@ def format_duration_delta(seconds: float) -> str:
     return f"+{format_duration(seconds)}"
 
 
-def format_prepare_breeze_timing(regression: dict) -> str | None:
-    """Format prepare breeze timing details for a job regression."""
-    if prepare_breeze := regression.get("prepare_breeze"):
-        return (
-            f"{PREPARE_BREEZE_STEP_PREFIX}: "
-            f"{format_duration(prepare_breeze['baseline'])} → "
-            f"{format_duration(prepare_breeze['latest'])} "
-            f"({format_duration_delta(prepare_breeze['increase'])})"
-        )
-    return None
-
-
 # Wall-clock shorter than this almost always means a run that was cancelled,
 # skipped by selective checks, or never really executed the test matrix — not a
 # representative "main build". Such runs would corrupt the duration baseline.
@@ -263,6 +259,34 @@ def get_prepare_breeze_step_duration(job: dict) -> float | 
None:
     return None
 
 
+def calculate_work_duration(job_data: JobDuration) -> float:
+    """Return a job's wall-clock with the image-build (prepare breeze) step 
removed.
+
+    The image build occasionally balloons — a cache miss forces a full rebuild
+    (minutes → tens of minutes) — which would otherwise inflate the job's total
+    and flag an unrelated job as "slower". The duration trend should track the
+    actual test/work time, so image build is discounted from it and watched
+    separately by :func:`detect_image_build_regression`.
+    """
+    prepare_breeze = job_data["prepare_breeze_duration"] or 0.0
+    return max(job_data["duration"] - prepare_breeze, 0.0)
+
+
+def calculate_image_build_seconds(jobs: dict[str, JobDuration]) -> float | 
None:
+    """Return a representative image-build duration for a run.
+
+    The same CI image is prepared by every job, so the median prepare-breeze
+    duration across the run's jobs is a robust single figure for that run
+    (ignoring jobs where the step is absent). None when no job recorded it.
+    """
+    values = [
+        job["prepare_breeze_duration"] for job in jobs.values() if 
job["prepare_breeze_duration"] is not None
+    ]
+    if not values:
+        return None
+    return median(values)
+
+
 def get_run_jobs(repo: str, run_id: int) -> dict[str, JobDuration]:
     """Return a mapping of job name -> duration details for a single run.
 
@@ -328,32 +352,38 @@ def detect_regression(
     return None
 
 
+def fetch_run_jobs_map(repo: str, runs: list[dict]) -> dict[int, dict[str, 
JobDuration]]:
+    """Fetch per-job durations once for every run, keyed by run id.
+
+    Shared by the per-job trend, the overall-run image-build discount, and the
+    image-build persistence check so the jobs of each run are fetched only 
once.
+    """
+    return {run["id"]: get_run_jobs(repo, run["id"]) for run in runs}
+
+
 def analyze_jobs(
-    repo: str,
+    jobs_by_run_id: dict[int, dict[str, JobDuration]],
     latest_runs: list[dict],
     baseline_runs: list[dict],
     min_baseline_runs: int,
     rel_threshold: float,
     min_abs_increase_seconds: float,
 ) -> list[dict]:
-    """Fetch per-job durations and return the jobs whose latest duration 
regressed."""
+    """Return the jobs whose latest duration regressed, image-build time 
excluded.
+
+    The comparison uses :func:`calculate_work_duration` (wall-clock minus the 
image-build
+    step) on both sides, so an occasional image rebuild spike does not flag a 
job
+    that did not actually get slower.
+    """
     latest_job_durations: dict[str, list[float]] = {}
-    latest_prepare_breeze_durations: dict[str, list[float]] = {}
     for run in latest_runs:
-        for name, job_data in get_run_jobs(repo, run["id"]).items():
-            latest_job_durations.setdefault(name, 
[]).append(job_data["duration"])
-            prepare_breeze_duration = job_data["prepare_breeze_duration"]
-            if prepare_breeze_duration is not None:
-                latest_prepare_breeze_durations.setdefault(name, 
[]).append(prepare_breeze_duration)
+        for name, job_data in jobs_by_run_id.get(run["id"], {}).items():
+            latest_job_durations.setdefault(name, 
[]).append(calculate_work_duration(job_data))
 
     baseline_job_durations: dict[str, list[float]] = {}
-    baseline_prepare_breeze_durations: dict[str, list[float]] = {}
     for run in baseline_runs:
-        for name, job_data in get_run_jobs(repo, run["id"]).items():
-            baseline_job_durations.setdefault(name, 
[]).append(job_data["duration"])
-            prepare_breeze_duration = job_data["prepare_breeze_duration"]
-            if prepare_breeze_duration is not None:
-                baseline_prepare_breeze_durations.setdefault(name, 
[]).append(prepare_breeze_duration)
+        for name, job_data in jobs_by_run_id.get(run["id"], {}).items():
+            baseline_job_durations.setdefault(name, 
[]).append(calculate_work_duration(job_data))
 
     regressions: list[dict] = []
     for name, latest_values in latest_job_durations.items():
@@ -364,16 +394,6 @@ def analyze_jobs(
             latest_values, baseline_values, rel_threshold, 
min_abs_increase_seconds
         )
         if regression:
-            latest_prepare_breeze_values = 
latest_prepare_breeze_durations.get(name, [])
-            baseline_prepare_breeze_values = 
baseline_prepare_breeze_durations.get(name, [])
-            if latest_prepare_breeze_values and 
len(baseline_prepare_breeze_values) >= min_baseline_runs:
-                latest_prepare_breeze = median(latest_prepare_breeze_values)
-                baseline_prepare_breeze = 
median(baseline_prepare_breeze_values)
-                regression["prepare_breeze"] = {
-                    "latest": latest_prepare_breeze,
-                    "baseline": baseline_prepare_breeze,
-                    "increase": latest_prepare_breeze - 
baseline_prepare_breeze,
-                }
             regression["job"] = name
             regressions.append(regression)
 
@@ -381,12 +401,73 @@ def analyze_jobs(
     return regressions
 
 
+def detect_image_build_regression(
+    runs: list[dict],
+    jobs_by_run_id: dict[int, dict[str, JobDuration]],
+    latest_runs_count: int,
+    min_baseline_runs: int,
+    rel_threshold: float,
+    min_abs_increase_seconds: float,
+    persistence_days: float,
+) -> dict | None:
+    """Flag the image build only when it has been slow for longer than 
``persistence_days``.
+
+    Image-build time is discounted from every other trend precisely because it
+    spikes on a one-off cache miss. A genuinely slow image (a bad base image, a
+    heavier install) shows up as an elevation that *persists* across runs. So 
we
+    compare each run's representative image-build time 
(:func:`calculate_image_build_seconds`)
+    against a baseline drawn from the oldest runs in the window, then require 
an
+    unbroken streak of elevated runs — starting at the most recent — that spans
+    more than ``persistence_days``. A single slow night never alerts.
+    """
+    per_run: list[tuple[dict, float]] = []
+    for run in runs:
+        seconds = calculate_image_build_seconds(jobs_by_run_id.get(run["id"], 
{}))
+        if seconds is not None:
+            per_run.append((run, seconds))
+    if len(per_run) < latest_runs_count + min_baseline_runs:
+        return None
+
+    # Baseline from the oldest runs in the window: a recent multi-day spike 
must not
+    # contaminate the baseline it is being measured against.
+    baseline = median([seconds for _, seconds in per_run[-min_baseline_runs:]])
+    threshold = baseline * (1 + rel_threshold)
+
+    streak: list[tuple[dict, float]] = []
+    for run, seconds in per_run:  # newest first
+        if seconds > threshold and (seconds - baseline) >= 
min_abs_increase_seconds:
+            streak.append((run, seconds))
+        else:
+            break
+    if not streak:
+        return None
+
+    newest_dt = parse_iso(streak[0][0].get("created_at"))
+    oldest_dt = parse_iso(streak[-1][0].get("created_at"))
+    if newest_dt is None or oldest_dt is None:
+        return None
+    span_seconds = (newest_dt - oldest_dt).total_seconds()
+    if span_seconds < persistence_days * 86400:
+        return None
+
+    latest = median([seconds for _, seconds in streak])
+    return {
+        "latest": latest,
+        "baseline": baseline,
+        "increase": latest - baseline,
+        "rel_increase": (latest - baseline) / baseline if baseline > 0 else 
0.0,
+        "elevated_runs": len(streak),
+        "span_days": span_seconds / 86400,
+    }
+
+
 def format_slack_message(
     repo: str,
     workflow: str,
     branch: str,
     overall_regression: dict | None,
     job_regressions: list[dict],
+    image_build_regression: dict | None,
     recent_runs: list[dict],
     rel_threshold: float,
     channel: str,
@@ -405,13 +486,33 @@ def format_slack_message(
                     f"CI run times on *{escape_slack_mrkdwn(branch)}* "
                     f"(`{escape_slack_mrkdwn(workflow)}`) have risen above the 
recent trend "
                     f"(baseline = median of the preceding runs; threshold = "
-                    f"+{int(rel_threshold * 100)}%)."
+                    f"+{int(rel_threshold * 100)}%). Image-build time is 
excluded from these "
+                    f"durations and reported separately."
                 ),
             },
         },
         {"type": "divider"},
     ]
 
+    if image_build_regression:
+        blocks.append(
+            {
+                "type": "section",
+                "text": {
+                    "type": "mrkdwn",
+                    "text": (
+                        f"🐳 *CI image build slow for 
{image_build_regression['span_days']:.1f} days* "
+                        f"(across {image_build_regression['elevated_runs']} 
runs) — not a one-off "
+                        f"cache miss:\n"
+                        f"• {PREPARE_BREEZE_STEP_PREFIX}: "
+                        
f"{format_duration(image_build_regression['baseline'])} → "
+                        
f"*{format_duration(image_build_regression['latest'])}* "
+                        f"(+{round(image_build_regression['rel_increase'] * 
100, 1)}%)"
+                    ),
+                },
+            }
+        )
+
     if overall_regression:
         blocks.append(
             {
@@ -430,16 +531,13 @@ def format_slack_message(
         )
 
     if job_regressions:
-        lines = ["*Jobs that got slower:*"]
+        lines = ["*Jobs that got slower (image build excluded):*"]
         for reg in job_regressions[:15]:
-            line = (
+            lines.append(
                 f"• *{escape_slack_mrkdwn(reg['job'])}* — "
                 f"{format_duration(reg['baseline'])} → 
*{format_duration(reg['latest'])}* "
                 f"(+{round(reg['rel_increase'] * 100, 1)}%)"
             )
-            if prepare_breeze_timing := format_prepare_breeze_timing(reg):
-                line += f"\n  {escape_slack_mrkdwn(prepare_breeze_timing)}"
-            lines.append(line)
         text = "\n".join(lines)
         if len(text) > 2900:
             text = text[:2900] + "\n_...truncated_"
@@ -482,6 +580,8 @@ def format_slack_message(
         fallback_parts.append(f"overall 
+{round(overall_regression['rel_increase'] * 100, 1)}%")
     if job_regressions:
         fallback_parts.append(f"{len(job_regressions)} slower job(s)")
+    if image_build_regression:
+        fallback_parts.append(f"image build slow 
{image_build_regression['span_days']:.1f}d")
     fallback = f"CI Duration Trend Alert on {branch}: " + ", 
".join(fallback_parts)
 
     return {
@@ -496,6 +596,7 @@ def write_step_summary(
     branch: str,
     overall_regression: dict | None,
     job_regressions: list[dict],
+    image_build_regression: dict | None,
     recent_runs: list[dict],
     baseline_count: int,
 ) -> None:
@@ -507,10 +608,23 @@ def write_step_summary(
     lines = [
         "## ⏱️ CI Duration Trend",
         "",
-        f"Workflow `{workflow}` on `{branch}` — baseline from {baseline_count} 
preceding runs.",
+        f"Workflow `{workflow}` on `{branch}` — baseline from {baseline_count} 
preceding runs. "
+        "Image-build time is excluded from run/job durations and tracked 
separately.",
         "",
     ]
 
+    if image_build_regression:
+        lines += [
+            f"### 🐳 CI image build slow for 
{image_build_regression['span_days']:.1f} days",
+            "",
+            f"- {PREPARE_BREEZE_STEP_PREFIX}: "
+            f"**{format_duration(image_build_regression['latest'])}** "
+            f"(baseline {format_duration(image_build_regression['baseline'])}, 
"
+            f"+{round(image_build_regression['rel_increase'] * 100, 1)}%) "
+            f"across {image_build_regression['elevated_runs']} runs",
+            "",
+        ]
+
     if overall_regression:
         lines += [
             "### ⚠️ Overall run regressed",
@@ -532,11 +646,8 @@ def write_step_summary(
             "|-----|----------|--------|----------|",
         ]
         for reg in job_regressions[:25]:
-            job = reg["job"]
-            if prepare_breeze_timing := format_prepare_breeze_timing(reg):
-                job += f"<br>{prepare_breeze_timing}"
             lines.append(
-                f"| {job} | {format_duration(reg['baseline'])} | "
+                f"| {reg['job']} | {format_duration(reg['baseline'])} | "
                 f"{format_duration(reg['latest'])} | 
+{round(reg['rel_increase'] * 100, 1)}% |"
             )
         lines.append("")
@@ -572,6 +683,7 @@ def main() -> None:
     rel_threshold = env_float("REL_THRESHOLD", 0.25)
     min_abs_increase_seconds = env_float("MIN_ABS_INCREASE_MINUTES", 5.0) * 60
     job_min_abs_increase_seconds = env_float("JOB_MIN_ABS_INCREASE_MINUTES", 
3.0) * 60
+    image_build_persistence_days = env_float("IMAGE_BUILD_PERSISTENCE_DAYS", 
2.0)
     do_analyze_jobs = env_bool("ANALYZE_JOBS", True)
     only_successful = env_bool("ONLY_SUCCESSFUL", True)
     channel = os.environ.get("SLACK_CHANNEL", "internal-airflow-ci-cd")
@@ -587,22 +699,32 @@ def main() -> None:
             f"Not enough runs to establish a trend "
             f"(found {len(runs)}, need {latest_runs_count + 
min_baseline_runs}). Skipping."
         )
-        _write_outputs(False, False, 0)
+        _write_outputs(False, False, 0, False)
         sys.exit(0)
 
     latest_runs = runs[:latest_runs_count]
     baseline_runs = runs[latest_runs_count:]
     print(f"Latest runs: {len(latest_runs)}; baseline runs: 
{len(baseline_runs)}.")
 
+    # Fetch each run's jobs once: they feed the image-build discount below, the
+    # per-job trend, and the image-build persistence check.
+    jobs_by_run_id = fetch_run_jobs_map(repo, runs) if do_analyze_jobs else {}
+
+    def calculate_adjusted_run_duration(run: dict) -> float:
+        """Run wall-clock with the image-build component of its critical path 
removed."""
+        image_build = 
calculate_image_build_seconds(jobs_by_run_id.get(run["id"], {})) or 0.0
+        return max(run["duration"] - image_build, 0.0)
+
     overall_regression = detect_regression(
-        [r["duration"] for r in latest_runs],
-        [r["duration"] for r in baseline_runs],
+        [calculate_adjusted_run_duration(r) for r in latest_runs],
+        [calculate_adjusted_run_duration(r) for r in baseline_runs],
         rel_threshold,
         min_abs_increase_seconds,
     )
     if overall_regression:
         print(
-            f"Overall regression: 
{format_duration(overall_regression['baseline'])} -> "
+            f"Overall regression (image build excluded): "
+            f"{format_duration(overall_regression['baseline'])} -> "
             f"{format_duration(overall_regression['latest'])} "
             f"(+{round(overall_regression['rel_increase'] * 100, 1)}%)"
         )
@@ -610,9 +732,10 @@ def main() -> None:
         print("Overall run duration is within the recent trend.")
 
     job_regressions: list[dict] = []
+    image_build_regression: dict | None = None
     if do_analyze_jobs:
         job_regressions = analyze_jobs(
-            repo,
+            jobs_by_run_id,
             latest_runs,
             baseline_runs,
             min_baseline_runs,
@@ -621,22 +744,60 @@ def main() -> None:
         )
         print(f"Jobs that regressed: {len(job_regressions)}")
 
-    has_regression = bool(overall_regression) or bool(job_regressions)
+        image_build_regression = detect_image_build_regression(
+            runs,
+            jobs_by_run_id,
+            latest_runs_count,
+            min_baseline_runs,
+            rel_threshold,
+            job_min_abs_increase_seconds,
+            image_build_persistence_days,
+        )
+        if image_build_regression:
+            print(
+                f"Image build slow for 
{image_build_regression['span_days']:.1f} days: "
+                f"{format_duration(image_build_regression['baseline'])} -> "
+                f"{format_duration(image_build_regression['latest'])}"
+            )
+        else:
+            print("Image build within trend (or not slow long enough to 
report).")
+
+    has_regression = bool(overall_regression) or bool(job_regressions) or 
bool(image_build_regression)
 
     if has_regression:
         slack_message = format_slack_message(
-            repo, workflow, branch, overall_regression, job_regressions, runs, 
rel_threshold, channel
+            repo,
+            workflow,
+            branch,
+            overall_regression,
+            job_regressions,
+            image_build_regression,
+            runs,
+            rel_threshold,
+            channel,
         )
         output_file.write_text(json.dumps(slack_message, indent=2))
         print(f"Slack message written to: {output_file}")
     else:
         print("No regression detected; no Slack message written.")
 
-    write_step_summary(workflow, branch, overall_regression, job_regressions, 
runs, len(baseline_runs))
-    _write_outputs(has_regression, bool(overall_regression), 
len(job_regressions))
+    write_step_summary(
+        workflow,
+        branch,
+        overall_regression,
+        job_regressions,
+        image_build_regression,
+        runs,
+        len(baseline_runs),
+    )
+    _write_outputs(
+        has_regression, bool(overall_regression), len(job_regressions), 
bool(image_build_regression)
+    )
 
 
-def _write_outputs(has_regression: bool, overall_regression: bool, 
regressed_jobs: int) -> None:
+def _write_outputs(
+    has_regression: bool, overall_regression: bool, regressed_jobs: int, 
image_build_regression: bool
+) -> None:
     """Write GitHub Actions outputs used to gate the Slack-notify step."""
     github_output = os.environ.get("GITHUB_OUTPUT")
     if not github_output:
@@ -645,6 +806,7 @@ def _write_outputs(has_regression: bool, 
overall_regression: bool, regressed_job
         f.write(f"has-regression={str(has_regression).lower()}\n")
         f.write(f"overall-regression={str(overall_regression).lower()}\n")
         f.write(f"regressed-jobs={regressed_jobs}\n")
+        
f.write(f"image-build-regression={str(image_build_regression).lower()}\n")
 
 
 if __name__ == "__main__":
diff --git a/scripts/tests/ci/test_analyze_ci_job_durations.py 
b/scripts/tests/ci/test_analyze_ci_job_durations.py
index 6bac585b60c..efb930539b3 100644
--- a/scripts/tests/ci/test_analyze_ci_job_durations.py
+++ b/scripts/tests/ci/test_analyze_ci_job_durations.py
@@ -325,37 +325,195 @@ class TestGetRunJobs:
             assert durations_module.get_run_jobs("apache/airflow", 2) == {}
 
 
+class TestWorkDuration:
+    def test_subtracts_image_build(self, durations_module):
+        assert (
+            durations_module.calculate_work_duration({"duration": 1200, 
"prepare_breeze_duration": 300})
+            == 900
+        )
+
+    def test_full_duration_when_no_image_build_step(self, durations_module):
+        assert (
+            durations_module.calculate_work_duration({"duration": 1200, 
"prepare_breeze_duration": None})
+            == 1200
+        )
+
+    def test_never_negative(self, durations_module):
+        assert (
+            durations_module.calculate_work_duration({"duration": 100, 
"prepare_breeze_duration": 300}) == 0
+        )
+
+
+class TestRunImageBuildSeconds:
+    def test_median_across_jobs(self, durations_module):
+        jobs = {
+            "a": {"duration": 0, "prepare_breeze_duration": 300},
+            "b": {"duration": 0, "prepare_breeze_duration": 500},
+            "c": {"duration": 0, "prepare_breeze_duration": 400},
+        }
+        assert durations_module.calculate_image_build_seconds(jobs) == 400
+
+    def test_ignores_jobs_without_the_step(self, durations_module):
+        jobs = {
+            "a": {"duration": 0, "prepare_breeze_duration": None},
+            "b": {"duration": 0, "prepare_breeze_duration": 500},
+        }
+        assert durations_module.calculate_image_build_seconds(jobs) == 500
+
+    def test_none_when_no_job_recorded_the_step(self, durations_module):
+        jobs = {"a": {"duration": 0, "prepare_breeze_duration": None}}
+        assert durations_module.calculate_image_build_seconds(jobs) is None
+
+
 class TestAnalyzeJobs:
     def test_reports_only_regressed_jobs_with_enough_baseline(self, 
durations_module):
         latest_runs = [{"id": 100}]
         baseline_runs = [{"id": i} for i in range(5)]
-
-        def fake_jobs(_repo, run_id):
-            if run_id == 100:
-                return {
-                    "slow-job": {"duration": 2700, "prepare_breeze_duration": 
900},
-                    "stable-job": {"duration": 600, "prepare_breeze_duration": 
None},
-                    "new-job": {"duration": 999, "prepare_breeze_duration": 
300},
-                }
-            # baseline runs
-            return {
+        # slow-job work time (image build excluded): latest 2400 vs baseline 
1500 -> +60%.
+        jobs_by_run_id = {
+            100: {
+                "slow-job": {"duration": 2700, "prepare_breeze_duration": 300},
+                "stable-job": {"duration": 600, "prepare_breeze_duration": 
None},
+                "new-job": {"duration": 999, "prepare_breeze_duration": 300},
+            }
+        }
+        for i in range(5):
+            jobs_by_run_id[i] = {
                 "slow-job": {"duration": 1800, "prepare_breeze_duration": 300},
                 "stable-job": {"duration": 590, "prepare_breeze_duration": 
None},
             }
 
-        with patch.object(durations_module, "get_run_jobs", 
side_effect=fake_jobs):
-            regressions = durations_module.analyze_jobs(
-                "apache/airflow",
-                latest_runs,
-                baseline_runs,
-                min_baseline_runs=5,
-                rel_threshold=0.25,
-                min_abs_increase_seconds=180,
-            )
+        regressions = durations_module.analyze_jobs(
+            jobs_by_run_id,
+            latest_runs,
+            baseline_runs,
+            min_baseline_runs=5,
+            rel_threshold=0.25,
+            min_abs_increase_seconds=180,
+        )
         names = [r["job"] for r in regressions]
         # slow-job regressed; stable-job did not; new-job lacks baseline 
samples
         assert names == ["slow-job"]
-        assert regressions[0]["prepare_breeze"] == {"latest": 900, "baseline": 
300, "increase": 600}
+        # Job regressions no longer carry image-build detail; that is reported 
separately.
+        assert "prepare_breeze" not in regressions[0]
+
+    def test_image_build_spike_alone_does_not_flag_a_job(self, 
durations_module):
+        """A job whose total ballooned only because the image build spiked is 
not flagged."""
+        latest_runs = [{"id": 100}]
+        baseline_runs = [{"id": i} for i in range(5)]
+        # latest total 1500 = +150% vs baseline 600, but work time (300) is 
unchanged.
+        jobs_by_run_id = {100: {"job": {"duration": 1500, 
"prepare_breeze_duration": 1200}}}
+        for i in range(5):
+            jobs_by_run_id[i] = {"job": {"duration": 600, 
"prepare_breeze_duration": 300}}
+
+        regressions = durations_module.analyze_jobs(
+            jobs_by_run_id,
+            latest_runs,
+            baseline_runs,
+            min_baseline_runs=5,
+            rel_threshold=0.25,
+            min_abs_increase_seconds=60,
+        )
+        assert regressions == []
+
+
+class TestDetectImageBuildRegression:
+    @staticmethod
+    def _runs_and_jobs(specs):
+        """Build (runs, jobs_by_run_id) from newest-first (run_id, created_at, 
image_seconds)."""
+        runs = []
+        jobs_by_run_id = {}
+        for run_id, created_at, image_seconds in specs:
+            runs.append({"id": run_id, "created_at": created_at})
+            jobs_by_run_id[run_id] = {
+                "job": {"duration": image_seconds, "prepare_breeze_duration": 
image_seconds}
+            }
+        return runs, jobs_by_run_id
+
+    def _call(self, durations_module, specs, persistence_days=2.0, 
min_abs=180.0):
+        runs, jobs_by_run_id = self._runs_and_jobs(specs)
+        return durations_module.detect_image_build_regression(
+            runs,
+            jobs_by_run_id,
+            latest_runs_count=1,
+            min_baseline_runs=3,
+            rel_threshold=0.25,
+            min_abs_increase_seconds=min_abs,
+            persistence_days=persistence_days,
+        )
+
+    def test_flags_when_slow_for_more_than_persistence_window(self, 
durations_module):
+        # Newest 4 runs (06-18..06-15, a 3-day span) elevated; older 4 at 
baseline.
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 1200),
+            (17, "2026-06-17T03:00:00Z", 1200),
+            (16, "2026-06-16T03:00:00Z", 1200),
+            (15, "2026-06-15T03:00:00Z", 1200),
+            (14, "2026-06-14T03:00:00Z", 300),
+            (13, "2026-06-13T03:00:00Z", 300),
+            (12, "2026-06-12T03:00:00Z", 300),
+            (11, "2026-06-11T03:00:00Z", 300),
+        ]
+        regression = self._call(durations_module, specs)
+        assert regression is not None
+        assert regression["baseline"] == 300
+        assert regression["latest"] == 1200
+        assert regression["elevated_runs"] == 4
+        assert round(regression["span_days"]) == 3
+
+    def test_ignores_one_off_spike(self, durations_module):
+        # Only the newest run is slow -> span 0 days -> below the 2-day window.
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 1200),
+            (17, "2026-06-17T03:00:00Z", 300),
+            (16, "2026-06-16T03:00:00Z", 300),
+            (15, "2026-06-15T03:00:00Z", 300),
+            (14, "2026-06-14T03:00:00Z", 300),
+        ]
+        assert self._call(durations_module, specs) is None
+
+    def test_non_consecutive_elevation_does_not_count(self, durations_module):
+        # A gap right after the newest run breaks the streak, so the span is 0.
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 1200),
+            (17, "2026-06-17T03:00:00Z", 300),
+            (16, "2026-06-16T03:00:00Z", 1200),
+            (15, "2026-06-15T03:00:00Z", 1200),
+            (14, "2026-06-14T03:00:00Z", 300),
+            (13, "2026-06-13T03:00:00Z", 300),
+        ]
+        assert self._call(durations_module, specs) is None
+
+    def test_no_alert_when_current_run_not_elevated(self, durations_module):
+        # Older runs were slow but the latest recovered -> no ongoing problem.
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 300),
+            (17, "2026-06-17T03:00:00Z", 1200),
+            (16, "2026-06-16T03:00:00Z", 1200),
+            (15, "2026-06-15T03:00:00Z", 1200),
+            (14, "2026-06-14T03:00:00Z", 300),
+            (13, "2026-06-13T03:00:00Z", 300),
+        ]
+        assert self._call(durations_module, specs) is None
+
+    def test_respects_absolute_floor(self, durations_module):
+        # Sustained but tiny elevation (300 -> 380) stays under the absolute 
floor.
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 380),
+            (17, "2026-06-17T03:00:00Z", 380),
+            (16, "2026-06-16T03:00:00Z", 380),
+            (15, "2026-06-15T03:00:00Z", 300),
+            (14, "2026-06-14T03:00:00Z", 300),
+            (13, "2026-06-13T03:00:00Z", 300),
+        ]
+        assert self._call(durations_module, specs, min_abs=180.0) is None
+
+    def test_none_when_not_enough_runs_with_image_data(self, durations_module):
+        specs = [
+            (18, "2026-06-18T03:00:00Z", 1200),
+            (17, "2026-06-17T03:00:00Z", 1200),
+        ]
+        assert self._call(durations_module, specs) is None
 
 
 class TestFormatSlackMessage:
@@ -373,6 +531,7 @@ class TestFormatSlackMessage:
             job_regressions=[
                 {"job": "Tests", "latest": 1500, "baseline": 1000, "increase": 
500, "rel_increase": 0.5}
             ],
+            image_build_regression=None,
             recent_runs=[{"run_number": 102, "html_url": "https://example/2";, 
"duration": 2700}],
             rel_threshold=0.25,
             channel="internal-airflow-ci-cd",
@@ -383,31 +542,32 @@ class TestFormatSlackMessage:
         assert "Tests" in text_blob
         assert "main" in msg["text"]
 
-    def test_includes_prepare_breeze_timing_when_available(self, 
durations_module):
+    def test_includes_image_build_section_when_present(self, durations_module):
         msg = durations_module.format_slack_message(
             repo="apache/airflow",
             workflow="ci-amd.yml",
             branch="main",
             overall_regression=None,
-            job_regressions=[
-                {
-                    "job": "Tests",
-                    "latest": 1500,
-                    "baseline": 1000,
-                    "increase": 500,
-                    "rel_increase": 0.5,
-                    "prepare_breeze": {"latest": 600, "baseline": 300, 
"increase": 300},
-                }
-            ],
+            job_regressions=[],
+            image_build_regression={
+                "latest": 1080,
+                "baseline": 300,
+                "increase": 780,
+                "rel_increase": 2.6,
+                "elevated_runs": 5,
+                "span_days": 3.0,
+            },
             recent_runs=[{"run_number": 102, "html_url": "https://example/2";, 
"duration": 2700}],
             rel_threshold=0.25,
             channel="internal-airflow-ci-cd",
         )
-        text_blob = json.dumps(msg)
-        assert "Prepare breeze &amp; CI image: 5m 00s" in text_blob
-        assert "10m 00s" in text_blob
+        text_blob = json.dumps(msg).lower()
+        # Present even though no job/overall regression: the image build alone 
triggers it.
+        assert "image build slow for 3.0 days" in text_blob
+        assert "18m 00s" in json.dumps(msg)  # 1080s latest
+        assert "image build slow" in msg["text"].lower()
 
-    def test_omits_prepare_breeze_timing_when_unavailable(self, 
durations_module):
+    def test_omits_image_build_section_when_none(self, durations_module):
         msg = durations_module.format_slack_message(
             repo="apache/airflow",
             workflow="ci-amd.yml",
@@ -416,8 +576,9 @@ class TestFormatSlackMessage:
             job_regressions=[
                 {"job": "Tests", "latest": 1500, "baseline": 1000, "increase": 
500, "rel_increase": 0.5}
             ],
+            image_build_regression=None,
             recent_runs=[{"run_number": 102, "html_url": "https://example/2";, 
"duration": 2700}],
             rel_threshold=0.25,
             channel="internal-airflow-ci-cd",
         )
-        assert "Prepare breeze" not in json.dumps(msg)
+        assert "image build slow" not in json.dumps(msg).lower()


Reply via email to