File size: 29,992 Bytes
2b2ba56
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
#!/usr/bin/env python3
"""Live-updating CI review comment.

Polls the GitHub Actions API for job statuses in the CI run, assembles
the review comment from whatever results are available, and upserts it as a
PR comment. Repeats every ``--interval`` seconds until all jobs are
completed (or ``--timeout`` is reached), so the comment updates in real time
as each job finishes.

The comment is identified by the ``<!-- hermes-ci-review-bot -->`` marker
β€” the same one ``assemble_review_comment.py`` uses β€” so it replaces any
previous comment from an earlier run.

This runs from ``.github/workflows/ci-review-comment.yml``, a separate
``workflow_run`` workflow. Thus ``CI_RUN_ID`` names the CI run to report
on, not the run that contains this script. (The variable cannot be
called ``GITHUB_RUN_ID``: the Actions runner sets the ``GITHUB_*``
defaults itself and ignores an ``env:`` override, so that name would
silently resolve to the poller's own run β€” which stays ``in_progress``
for as long as the poller runs, deadlocking it against itself.)
The poller reports on runs that
it does not belong to. This is also how it covers a workflow that CI does
not contain: ``WATCH_WORKFLOWS`` names sibling workflows that the same
commit triggered (the Docker image build). Their jobs join the comment.

Architecture:

  - :func:`classify_jobs` (pure, testable) β€” takes a list of raw API job
    dicts and returns ``(completed, pending, job_urls)`` where ``completed``
    is a ``{name: result}`` dict (for :func:`assemble_review_comment.assemble`)
    and ``pending`` is a list of job names still running.

  - :func:`select_watched_runs` (pure, testable) β€” picks the sibling runs
    to merge in, newest attempt per workflow.

  - :func:`find_comment_id` / :func:`upsert_comment` β€” thin API wrappers.

  - :func:`fetch_all_review_statuses` β€” lists all ``review-status-*``
    artifacts on the CI run (GitHub attaches reusable-workflow
    artifacts to the caller run), downloads each, parses the
    ``review_status=`` line from ``review-status.json``, and merges into
    one array. Recomputed from source every poll cycle, so statuses
    appear as soon as each job uploads its artifact.

  - :func:`run` β€” the polling loop. Calls the API, classifies,
    fetches artifacts, assembles, upserts, sleeps, repeats. Before
    its final exit, it gives downstream jobs a short grace period
    to appear.

The orchestrator job names (detect, all-checks-pass, comment-live, etc.)
are excluded from the comment β€” they're infrastructure, not review signal.
"""

from __future__ import annotations

import argparse
import json
import os
import shutil
import sys
import time
import urllib.error
import urllib.request
import zipfile
from pathlib import Path

API_BASE = "https://api.github.com"

# Job names that are infrastructure (this script, the gate, the detector)
# and should never appear in the review comment.
_INFRA_JOBS = frozenset({
    "detect",
    "all-checks-pass",
    "comment-pending",
    "comment-results",
    "comment-live",
    "CI review comment (pending)",
    "CI review comment (results)",
    "CI review comment (live)",
    "All required checks pass",
    "Detect affected areas",
})

# Map GitHub API conclusion values to our result strings.
_CONCLUSION_MAP = {
    "success": "success",
    "failure": "failure",
    "skipped": "skipped",
    "cancelled": "skipped",
    "neutral": "skipped",
    "timed_out": "failure",
    "action_required": "skipped",
}

def classify_jobs(api_jobs: list[dict]) -> tuple[dict[str, str], list[str], dict[str, str]]:
    """Classify raw API job dicts into completed + pending + job_urls.

    Returns ``(completed, pending, job_urls)``:

    - ``completed``: ``{job_name: result}`` where result is
      ``"success"`` / ``"failure"`` / ``"skipped"``. Only non-infra jobs
      that have finished.
    - ``pending``: list of job names still running (in_progress / queued
      / waiting). Excludes infra jobs.
    - ``job_urls``: ``{job_name: html_url}`` β€” direct links to each
      job's logs page, for the assembler to use in ❌ Error links.

    The API returns orchestrator-level jobs and sub-workflow jobs
    (workflow_call) in separate runs β€” :func:`collect_run_jobs` merges
    them. Each sub-workflow job has a ``_workflow_name`` prefix so the
    display name is ``"Workflow / job"``.
    """
    completed: dict[str, str] = {}
    pending: list[str] = []
    job_urls: dict[str, str] = {}

    for job in api_jobs:
        name = job.get("name", "unknown")
        if job.get("_workflow_name"):
            name = f"{job['_workflow_name']} / {name}"
        if name in _INFRA_JOBS:
            continue
        status = job.get("status", "")
        conclusion = job.get("conclusion", "")
        html_url = job.get("html_url", "")

        if html_url:
            job_urls[name] = html_url

        if status in ("in_progress", "queued", "waiting"):
            pending.append(name)
        elif status == "completed":
            result = _CONCLUSION_MAP.get(conclusion, "skipped")
            completed[name] = result
        # else: unknown status β†’ skip

    return completed, pending, job_urls


# ---------------------------------------------------------------------------
# API helpers
# ---------------------------------------------------------------------------


def _api_request(url: str, token: str) -> dict:
    """Authenticated GitHub API GET (single page)."""
    req = urllib.request.Request(url, headers={
        "Authorization": f"Bearer {token}",
        "Accept": "application/vnd.github+json",
        "X-GitHub-Api-Version": "2022-11-28",
        "User-Agent": "ci-live-comment",
    })
    with urllib.request.urlopen(req) as resp:
        data: dict = json.loads(resp.read())
        return data


def _api_get_paginated(url: str, token: str, list_key: str | None = None) -> list:
    """Authenticated GitHub API GET with pagination."""
    results: list = []
    while url:
        req = urllib.request.Request(url, headers={
            "Authorization": f"Bearer {token}",
            "Accept": "application/vnd.github+json",
            "X-GitHub-Api-Version": "2022-11-28",
            "User-Agent": "ci-live-comment",
        })
        with urllib.request.urlopen(req) as resp:
            data = json.loads(resp.read())
            link_header = resp.headers.get("Link", "")

        if list_key:
            results.extend(data.get(list_key, []))
        elif isinstance(data, list):
            results.extend(data)
        else:
            return data

        next_url = None
        for part in link_header.split(","):
            part = part.strip()
            if 'rel="next"' in part:
                next_url = part[part.find("<") + 1:part.find(">")]
                break
        url = next_url

    return results


def select_watched_runs(
    runs: list[dict], watch_names: list[str], exclude_run_id: str = "",
) -> list[dict]:
    """Pick the sibling runs whose jobs belong in the comment.

    ``runs`` is the API's run list for one commit. ``watch_names`` holds
    workflow names from ``WATCH_WORKFLOWS``. One commit can have more than
    one run of the same workflow, after a rerun or a new push. Thus this
    keeps only the newest run for each workflow name. An older attempt
    reports results that a rerun replaced.

    ``exclude_run_id`` removes the CI run itself when its name is also in
    ``watch_names``.
    """
    newest: dict[str, dict] = {}
    wanted = {n.strip() for n in watch_names if n.strip()}

    for candidate in runs:
        name = str(candidate.get("name", ""))
        if name not in wanted:
            continue
        if exclude_run_id and str(candidate.get("id", "")) == str(exclude_run_id):
            continue
        current = newest.get(name)
        if current is None or str(candidate.get("created_at", "")) > str(current.get("created_at", "")):
            newest[name] = candidate

    return list(newest.values())


def runs_all_completed(runs: list[dict]) -> bool:
    """True only when every run in the list reports ``status: completed``.

    The job list alone cannot answer "is CI done": a run that GitHub just
    created has no jobs yet, and a mid-run poll can catch the moment where
    every visible job finished but a downstream sub-workflow has not
    spawned its jobs. Both look identical to "all done" at the job level.
    The run's own ``status`` is the authoritative signal, so the poller
    must not exit while any relevant run is still ``queued`` or
    ``in_progress``. An empty list is not done β€” it means the poller has
    no run information at all.
    """
    return bool(runs) and all(str(r.get("status", "")) == "completed" for r in runs)


def collect_run_jobs(
    token: str, repo: str, run_id: str, watch_workflows: list[str] | None = None,
) -> tuple[list[dict], bool]:
    """Collect all jobs in the CI run + any watched sibling runs.

    Returns ``(jobs, runs_completed)``: a flat list of job dicts (same
    shape as the API returns, plus ``_workflow_name`` on jobs from a
    watched run), and whether the CI run and every selected watched run
    report ``status: completed`` (see :func:`runs_all_completed`).

    Reusable-workflow (``workflow_call``) jobs need no special handling:
    GitHub flattens them into the caller run's job list, already named
    ``\"Workflow / job\"``. Watched runs are separate top-level runs
    (the Docker image build), so their jobs are fetched per run and
    prefixed here.
    """
    owner, repo_name = repo.split("/")
    run_info = _api_request(f"{API_BASE}/repos/{owner}/{repo_name}/actions/runs/{run_id}", token)
    head_sha = run_info.get("head_sha", "")

    # CI run jobs (includes every reusable-workflow job).
    all_jobs: list[dict] = []
    orch_jobs = _api_get_paginated(
        f"{API_BASE}/repos/{owner}/{repo_name}/actions/runs/{run_id}/jobs",
        token, list_key="jobs",
    )

    # Skip workflow-call placeholder steps (they're sub-workflow triggers,
    # not review signal), but KEEP in_progress / queued jobs so the poller
    # knows they're still running.
    for job in orch_jobs:
        steps = job.get("steps") or []
        if any(s.get("name", "").startswith("Run ./.github/workflows/") for s in steps):
            continue
        all_jobs.append(job)

    if not watch_workflows or not head_sha:
        return all_jobs, runs_all_completed([run_info])

    # Watched sibling runs for the same commit. A run can be absent on the
    # first polls. Then classify_jobs() shows nothing for it.
    sibling_runs = _api_get_paginated(
        f"{API_BASE}/repos/{owner}/{repo_name}/actions/runs?head_sha={head_sha}&per_page=100",
        token, list_key="workflow_runs",
    )
    relevant_runs = [run_info]
    for watched in select_watched_runs(sibling_runs, watch_workflows, exclude_run_id=run_id):
        relevant_runs.append(watched)
        watched_jobs = _api_get_paginated(
            f"{API_BASE}/repos/{owner}/{repo_name}/actions/runs/{watched['id']}/jobs",
            token, list_key="jobs",
        )
        for job in watched_jobs:
            job["_workflow_name"] = watched.get("name", "")
            all_jobs.append(job)

    return all_jobs, runs_all_completed(relevant_runs)


def find_comment_id(token: str, repo: str, pr_number: str) -> int | None:
    """Find our existing review comment by marker prefix."""
    owner, repo_name = repo.split("/")
    comments = _api_get_paginated(
        f"{API_BASE}/repos/{owner}/{repo_name}/issues/{pr_number}/comments",
        token,
    )
    for c in comments:
        body = c.get("body", "") if isinstance(c, dict) else ""
        if body.startswith("<!-- hermes-ci-review-bot -->"):
            return c.get("id") if isinstance(c, dict) else None
    return None


def upsert_comment(
    token: str, repo: str, pr_number: str, body: str, comment_id: int | None = None
) -> int | None:
    """Create or update the review comment. Returns the comment ID."""
    owner, repo_name = repo.split("/")
    if comment_id is None:
        comment_id = find_comment_id(token, repo, pr_number)

    if comment_id:
        url = f"{API_BASE}/repos/{owner}/{repo_name}/issues/comments/{comment_id}"
        method = "PATCH"
    else:
        url = f"{API_BASE}/repos/{owner}/{repo_name}/issues/{pr_number}/comments"
        method = "POST"

    data = json.dumps({"body": body}).encode("utf-8")
    req = urllib.request.Request(url, data=data, method=method, headers={
        "Authorization": f"Bearer {token}",
        "Accept": "application/vnd.github+json",
        "X-GitHub-Api-Version": "2022-11-28",
        "Content-Type": "application/json",
        "User-Agent": "ci-live-comment",
    })
    try:
        with urllib.request.urlopen(req) as resp:
            result = json.loads(resp.read())
            return result.get("id")
    except urllib.error.HTTPError as e:
        print(f"  API error {e.code}: {e.reason}", file=sys.stderr)
        return None


# ---------------------------------------------------------------------------
# Artifact fetching (dynamic review-status artifacts)
# ---------------------------------------------------------------------------

# Prefix for all review-status artifacts uploaded by status-producing jobs.
# Each job uploads a ``review-status-<name>`` artifact containing a
# ``review-status.json`` file in GITHUB_OUTPUT format:
#   review_status=<json array of {source, results: [...]} objects>
_REVIEW_STATUS_ARTIFACT_PREFIX = "review-status-"


def _list_artifacts(token: str, repo: str, run_id: str) -> list[dict]:
    """List artifacts for a given run (paginated)."""
    owner, repo_name = repo.split("/")
    return _api_get_paginated(
        f"{API_BASE}/repos/{owner}/{repo_name}/actions/runs/{run_id}/artifacts",
        token, list_key="artifacts",
    )


class _NoRedirectHandler(urllib.request.HTTPRedirectHandler):
    """Redirect handler that never follows β€” used to capture the Location."""

    def redirect_request(self, *args, **kwargs):
        return None


def _download_artifact(
    token: str, repo: str, artifact: dict, dest_dir: Path,
) -> Path | None:
    """Download a single artifact zip via the API and extract it.

    Returns the path to ``review-status.json`` inside the extracted dir,
    or ``None`` if the download or extraction failed.
    """
    owner, repo_name = repo.split("/")
    archive_download_url = artifact.get("archive_download_url", "")
    if not archive_download_url:
        return None

    # The archive_download_url is an API URL that 302s to a signed blob
    # URL. Hop 1 authenticates to the API; hop 2 follows the redirect
    # WITHOUT the Authorization header β€” the blob rejects a request that
    # carries both a SAS token and an Authorization header (401).
    opener = urllib.request.build_opener(_NoRedirectHandler)
    location = ""
    try:
        opener.open(urllib.request.Request(archive_download_url, headers={
            "Authorization": f"Bearer {token}",
            "Accept": "application/vnd.github+json",
            "X-GitHub-Api-Version": "2022-11-28",
            "User-Agent": "ci-live-comment",
        }), timeout=30)
    except urllib.error.HTTPError as e:
        location = e.headers.get("Location", "") if e.code == 302 else ""
    except Exception:
        location = ""
    if not location:
        return None

    zip_path = dest_dir / f"{artifact['name']}.zip"
    try:
        # No auth headers here; further redirects are safe to follow.
        with urllib.request.urlopen(
            urllib.request.Request(location, headers={"User-Agent": "ci-live-comment"}),
            timeout=60,
        ) as resp:
            zip_path.write_bytes(resp.read())
    except Exception:
        return None

    extract_dir = dest_dir / artifact["name"]
    extract_dir.mkdir(parents=True, exist_ok=True)
    try:
        with zipfile.ZipFile(zip_path) as zf:
            if any(".." in name or name.startswith("/") for name in zf.namelist()):
                return None
            zf.extractall(extract_dir)
    except Exception:
        return None

    status_file = extract_dir / "review-status.json"
    return status_file if status_file.exists() else None


def _parse_status_file(status_file: Path) -> list[dict]:
    """Parse a review-status.json file in GITHUB_OUTPUT format."""
    try:
        content = status_file.read_text(encoding="utf-8").strip()
        if content.startswith("review_status="):
            content = content[len("review_status="):]
        statuses = json.loads(content)
        if isinstance(statuses, list):
            return statuses
    except (json.JSONDecodeError, OSError):
        pass
    return []


def fetch_all_review_statuses(
    token: str, repo: str, run_id: str,
) -> list[dict]:
    """Fetch and merge all review-status artifacts from the run.

    Lists artifacts with the ``review-status-`` prefix on the orchestrator
    run, downloads each, parses the ``review-status.json`` inside, and
    merges into a single flat array. GitHub attaches artifacts uploaded by
    reusable workflow jobs to the caller run, so one listing covers every
    status-producing job.

    Returns the merged list of ``{source, results: [...]}`` objects.
    Artifacts that don't exist yet or fail to parse are silently skipped.
    """
    all_statuses: list[dict] = []
    temp_base = Path("/tmp/review-status-artifacts")

    try:
        artifacts = _list_artifacts(token, repo, run_id)
    except Exception:
        return all_statuses

    rs_artifacts = [
        a for a in artifacts
        if a.get("name", "").startswith(_REVIEW_STATUS_ARTIFACT_PREFIX)
    ]
    if not rs_artifacts:
        return all_statuses

    # Clean temp dir for this run's artifacts.
    run_dl_dir = temp_base / str(run_id)
    if run_dl_dir.exists():
        shutil.rmtree(run_dl_dir)
    run_dl_dir.mkdir(parents=True, exist_ok=True)

    for artifact in rs_artifacts:
        status_file = _download_artifact(token, repo, artifact, run_dl_dir)
        if status_file is None:
            continue
        statuses = _parse_status_file(status_file)
        all_statuses.extend(statuses)

    # A re-run can leave several non-expired artifacts with the same name,
    # each carrying the same source β€” dedupe by source so the comment
    # doesn't render duplicate sections.
    seen: set[str] = set()
    deduped: list[dict] = []
    for status in all_statuses:
        src = status.get("source", "")
        if src in seen:
            continue
        if src:
            seen.add(src)
        deduped.append(status)
    return deduped


# ---------------------------------------------------------------------------
# Comment assembly
# ---------------------------------------------------------------------------


def _import_assembler():
    """Import assemble_review_comment.py from the same directory."""
    here = Path(__file__).resolve().parent
    sys.path.insert(0, str(here))
    import assemble_review_comment as asm
    return asm


def build_comment_body(
    asm_mod,
    completed: dict[str, str],
    pending: list[str],
    run_url: str,
    job_urls: dict[str, str],
    review_statuses_json: str,
    commit_info: str = "",
    waiting: bool = False,
) -> str:
    """Assemble the comment body from current job states + static inputs."""
    needs_json = json.dumps(completed) if completed else ""

    return asm_mod.assemble(
        needs_json=needs_json,
        run_url=run_url,
        job_urls=job_urls,
        review_statuses_json=review_statuses_json,
        pending_jobs=pending if pending else None,
        commit_info=commit_info,
        waiting=waiting,
    )


def _commit_info_for_state(commit_info: str, pending: bool) -> str:
    """Use past tense in the final comment after every CI job completes."""
    if pending:
        return commit_info
    return commit_info.replace("<sub>running on ", "<sub>ran on ", 1)


# ---------------------------------------------------------------------------
# Polling loop
# ---------------------------------------------------------------------------


def run(
    token: str,
    repo: str,
    run_id: str,
    pr_number: str,
    run_url: str,
    commit_info: str = "",
    interval: int = 15,
    timeout: int = 1800,
    dry_run: bool = False,
    watch_workflows: list[str] | None = None,
) -> int:
    """Poll for job statuses and update the PR comment until all done.

    Always returns 0. The poller reports on the CI run from a different run.
    Thus a failed CI job is not a failure of this job. The CI run has its
    own gate, which reports that. Comment posting is best-effort.
    """
    asm = _import_assembler()
    start = time.time()
    last_body = ""
    quiet_grace_used = False
    prev_completed: dict[str, str] = {}
    prev_pending: list[str] = []
    prev_artifact_count = 0

    while True:
        elapsed = time.time() - start
        if elapsed > timeout:
            print(f"Timeout ({timeout}s) reached β€” stopping poll.", file=sys.stderr)
            break

        try:
            jobs, runs_completed = collect_run_jobs(token, repo, run_id, watch_workflows)
        except Exception as e:
            print(f"  API error collecting jobs: {e}", file=sys.stderr)
            time.sleep(interval)
            continue

        completed, pending, job_urls = classify_jobs(jobs)
        total = len(completed) + len(pending)
        infra_count = len(jobs) - total
        print(f"  [{elapsed:.0f}s] fetched {len(jobs)} jobs from API "
              f"({infra_count} infra filtered) β†’ {len(completed)} completed, "
              f"{len(pending)} pending ({total} review jobs)")

        # Log transitions since last poll.
        new_completed = {k: v for k, v in completed.items() if k not in prev_completed}
        new_pending = [j for j in pending if j not in prev_pending]
        gone_pending = [j for j in prev_pending if j not in pending and j not in completed]
        if new_completed:
            parts = [f"{name}={result}" for name, result in new_completed.items()]
            print(f"  β†’ {len(new_completed)} job(s) newly completed: {', '.join(parts)}")
        if new_pending:
            print(f"  β†’ {len(new_pending)} job(s) newly appeared: {', '.join(new_pending)}")
        if gone_pending:
            print(f"  β†’ {len(gone_pending)} job(s) disappeared from pending: {', '.join(gone_pending)}")

        # Dynamically fetch all review-status artifacts from the run.
        artifact_statuses = fetch_all_review_statuses(token, repo, run_id)
        artifact_count_changed = len(artifact_statuses) != prev_artifact_count
        if artifact_count_changed:
            print(f"  Found {len(artifact_statuses)} review status entries from artifacts "
                  f"(was {prev_artifact_count} last poll)")
        prev_artifact_count = len(artifact_statuses)

        merged_json = json.dumps(artifact_statuses) if artifact_statuses else ""
        # The run status is authoritative for "done": an empty job list on
        # a run that is still queued/in_progress means GitHub has not
        # spawned the jobs yet, not that everything passed.
        all_done = not pending and runs_completed
        current_commit_info = _commit_info_for_state(commit_info, pending=not all_done)

        body = build_comment_body(
            asm, completed, pending, run_url, job_urls,
            merged_json,
            current_commit_info,
            waiting=not runs_completed,
        )

        if body != last_body:
            change_reasons = []
            if new_completed:
                change_reasons.append(f"{len(new_completed)} new completion(s)")
            if new_pending:
                change_reasons.append(f"{len(new_pending)} new pending job(s)")
            if gone_pending:
                change_reasons.append(f"{len(gone_pending)} job(s) left pending")
            if artifact_count_changed:
                change_reasons.append("artifact statuses updated")
            if not change_reasons:
                change_reasons.append("initial post")
            reason = "; ".join(change_reasons)

            if dry_run:
                print(f"  Comment body changed ({reason}) β€” DRY RUN:")
                print("--- DRY RUN β€” comment body ---")
                print(body)
                print("--- END ---")
            else:
                cid = upsert_comment(token, repo, pr_number, body)
                if cid:
                    print(f"  Updated comment {cid} ({reason})")
                else:
                    print(f"  Failed to update comment ({reason}, will retry)", file=sys.stderr)
            last_body = body
        else:
            if pending:
                print(f"  No change since last poll. Still waiting on: {', '.join(pending)}")
            else:
                print("  No change since last poll.")

        prev_completed = completed
        prev_pending = pending

        if all_done and not quiet_grace_used:
            quiet_grace_used = True
            print("  No jobs pending and runs report completed β€” "
                  "waiting 10s for downstream jobs to appear.")
            time.sleep(10)
            continue

        if all_done:
            failed = [name for name, result in completed.items() if result == "failure"]
            if failed:
                print(f"  All jobs done, {len(failed)} failed: {', '.join(failed)}")
            else:
                print("  All jobs completed β€” done.")
            break

        if not pending:
            print("  No visible jobs pending, but a run is still queued or "
                  "in progress β€” waiting for its jobs to appear.")

        quiet_grace_used = False
        time.sleep(interval)

    return 0


def parse_watch_workflows(raw: str) -> list[str]:
    """Parse the ``WATCH_WORKFLOWS`` value into workflow names.

    One name per line. Not comma-separated: a workflow name can contain a
    comma ("Docker Build, Test, and Publish").
    """
    return [name.strip() for name in raw.splitlines() if name.strip()]


def resolve_pr_number(token: str, repo: str, head_sha: str) -> str:
    """Find the PR number for a commit when the event payload has none.

    ``workflow_run.pull_requests`` is empty for some runs. The poller has no
    comment to post without a number.
    """
    if not head_sha:
        return ""
    owner, repo_name = repo.split("/")
    try:
        results = _api_get_paginated(
            f"{API_BASE}/repos/{owner}/{repo_name}/commits/{head_sha}/pulls",
            token,
        )
    except Exception as e:
        print(f"  API error resolving PR number: {e}", file=sys.stderr)
        return ""
    for item in results:
        if isinstance(item, dict) and item.get("state") == "open":
            return str(item.get("number", ""))
    return ""


def main() -> int:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--interval", type=int, default=15,
                        help="Seconds between polls (default: 15).")
    parser.add_argument("--timeout", type=int, default=1800,
                        help="Max seconds to poll before giving up (default: 1800).")
    parser.add_argument("--dry-run", action="store_true",
                        help="Print comment body instead of posting to PR.")
    args = parser.parse_args()

    token = os.environ.get("GITHUB_TOKEN", "")
    repo = os.environ.get("GITHUB_REPOSITORY", "")
    run_id = os.environ.get("CI_RUN_ID", "")
    pr_number = os.environ.get("PR_NUMBER", "")
    run_url = os.environ.get("RUN_URL", "")

    # Sibling workflows to merge into the comment, one name per line. Their
    # runs are separate from the CI run, so the poller resolves them by name.
    watch_workflows = parse_watch_workflows(os.environ.get("WATCH_WORKFLOWS", ""))

    if not args.dry_run:
        if not token:
            print("GITHUB_TOKEN is required", file=sys.stderr)
            return 1
        if not repo:
            print("GITHUB_REPOSITORY is required", file=sys.stderr)
            return 1
        if not run_id:
            print("CI_RUN_ID is required", file=sys.stderr)
            return 1

    # Build commit info line from env vars (set by ci-review-comment.yml).
    commit_sha = os.environ.get("COMMIT_SHA", "")
    commit_msg = os.environ.get("COMMIT_MESSAGE", "")

    if not pr_number and not args.dry_run:
        pr_number = resolve_pr_number(token, repo, commit_sha)
        if not pr_number:
            print("No PR number found β€” nothing to comment on.", file=sys.stderr)
            return 0
        print(f"Resolved PR #{pr_number} from commit {commit_sha[:7]}")

    commit_url = os.environ.get("COMMIT_URL", "")
    if not commit_url and commit_sha and pr_number:
        server = os.environ.get("GITHUB_SERVER_URL", "https://github.com")
        commit_url = f"{server}/{repo}/pull/{pr_number}/commits/{commit_sha}"

    commit_info = ""
    if commit_sha:
        short_sha = commit_sha[:7]
        if commit_msg:
            # Truncate commit message to first line, max 60 chars.
            first_line = commit_msg.split("\n")[0][:60]
            if commit_url:
                commit_info = f"<sub>running on [{short_sha}]({commit_url}) β€” {first_line}</sub>"
            else:
                commit_info = f"<sub>running on {short_sha} β€” {first_line}</sub>"
        elif commit_url:
            commit_info = f"<sub>running on [{short_sha}]({commit_url})</sub>"
        else:
            commit_info = f"<sub>running on {short_sha}</sub>"

    return run(
        token=token,
        repo=repo,
        run_id=run_id,
        pr_number=pr_number,
        run_url=run_url,
        commit_info=commit_info,
        interval=args.interval,
        timeout=args.timeout,
        dry_run=args.dry_run,
        watch_workflows=watch_workflows,
    )


if __name__ == "__main__":
    sys.exit(main())