kaxil commented on code in PR #71072:
URL: https://github.com/apache/airflow/pull/71072#discussion_r4150699842


##########
airflow-core/tests/unit/jobs/test_scheduler_job.py:
##########
@@ -13244,6 +13296,594 @@ def 
test_partition_cap_at_n_minus_one_leaves_one_pending(dag_maker: DagMaker, se
     assert partition_dags == {"cap-consumer-n-minus-one"}
 
 
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
[email protected](
+    ("cap", "partition_keys", "expect_cap_hit"),
+    [
+        pytest.param(2, ["k1", "k2", "k3"], True, id="over-cap"),
+        pytest.param(3, ["k1", "k2"], False, id="under-cap"),
+        # Regression guard for the ``backlog_total > cap`` boundary: 
``backlog_total == cap``
+        # must not be misreported as a backlog (that would be the "exactly 
cap, nothing left"
+        # case wrongly reported as "more than cap, backlog remains").
+        pytest.param(3, ["k1", "k2", "k3"], False, id="exactly-cap"),
+    ],
+)
+def test_partition_cap_reporting(
+    dag_maker: DagMaker,
+    session: Session,
+    caplog,
+    cap: int,
+    partition_keys: list[str],
+    expect_cap_hit: bool,
+):
+    """Only a pending count strictly above the cap reports a backlog, via log 
and audit `Log` row."""
+    suffix = f"{cap}-{len(partition_keys)}"
+    _make_n_satisfied_apdrs(
+        consumer_dag_id=f"cap-consumer-{suffix}",
+        asset=Asset(name=f"asset-cap-{suffix}"),
+        partition_keys=partition_keys,
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    dag_id = f"cap-consumer-{suffix}"
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    if expect_cap_hit:
+        assert {
+            "event": "Reached the per-tick cap on pending partitioned Dag 
runs; the remaining backlog "
+            "will be evaluated over subsequent scheduler ticks",
+            "cap": cap,
+            "backlog_total": 3,
+            "dag_ids": [dag_id],
+        } in caplog
+        assert [row.event for row in audit_rows] == ["partition Dag run cap 
reached"]
+        extra = audit_rows[0].extra
+        assert extra is not None
+        assert f"Affected dag_ids (highest pending count first): {dag_id}." in 
extra
+    else:
+        assert "Reached the per-tick cap on pending partitioned Dag runs" not 
in caplog

Review Comment:
   `StructlogCapture.__contains__` compares `e["event"] == target` exactly, and 
the real event continues with "; the remaining backlog will be evaluated over 
subsequent scheduler ticks". So this `not in` always passes, and the same goes 
for L13394 and L13449. Pulling the full message into a module constant and 
using it here (and everywhere else the string is repeated) would fix both 
problems.



##########
airflow-core/tests/unit/jobs/test_scheduler_job.py:
##########
@@ -13244,6 +13296,594 @@ def 
test_partition_cap_at_n_minus_one_leaves_one_pending(dag_maker: DagMaker, se
     assert partition_dags == {"cap-consumer-n-minus-one"}
 
 
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
[email protected](
+    ("cap", "partition_keys", "expect_cap_hit"),
+    [
+        pytest.param(2, ["k1", "k2", "k3"], True, id="over-cap"),
+        pytest.param(3, ["k1", "k2"], False, id="under-cap"),
+        # Regression guard for the ``backlog_total > cap`` boundary: 
``backlog_total == cap``
+        # must not be misreported as a backlog (that would be the "exactly 
cap, nothing left"
+        # case wrongly reported as "more than cap, backlog remains").
+        pytest.param(3, ["k1", "k2", "k3"], False, id="exactly-cap"),
+    ],
+)
+def test_partition_cap_reporting(
+    dag_maker: DagMaker,
+    session: Session,
+    caplog,
+    cap: int,
+    partition_keys: list[str],
+    expect_cap_hit: bool,
+):
+    """Only a pending count strictly above the cap reports a backlog, via log 
and audit `Log` row."""
+    suffix = f"{cap}-{len(partition_keys)}"
+    _make_n_satisfied_apdrs(
+        consumer_dag_id=f"cap-consumer-{suffix}",
+        asset=Asset(name=f"asset-cap-{suffix}"),
+        partition_keys=partition_keys,
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    dag_id = f"cap-consumer-{suffix}"
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    if expect_cap_hit:
+        assert {
+            "event": "Reached the per-tick cap on pending partitioned Dag 
runs; the remaining backlog "
+            "will be evaluated over subsequent scheduler ticks",
+            "cap": cap,
+            "backlog_total": 3,
+            "dag_ids": [dag_id],
+        } in caplog
+        assert [row.event for row in audit_rows] == ["partition Dag run cap 
reached"]
+        extra = audit_rows[0].extra
+        assert extra is not None
+        assert f"Affected dag_ids (highest pending count first): {dag_id}." in 
extra
+    else:
+        assert "Reached the per-tick cap on pending partitioned Dag runs" not 
in caplog
+        assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_already_fired_decoy(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    ``exactly-cap`` in :func:`test_partition_cap_reporting` cannot catch a 
filter drift in the
+    count query (e.g. a dropped ``created_dag_run_id.is_(None)`` filter): 
every row in that
+    table at that point is a genuine, unfired, non-stale APDR, so a looser 
filter has nothing
+    extra to over-count. This test adds a decoy APDR that already fired
+    (``created_dag_run_id`` set) — a naive/unfiltered ``count()`` would 
wrongly include it.
+    """
+    cap = 3
+    asset = Asset(name="asset-cap-decoy")
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-decoy",
+        asset=asset,
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+    decoy = _produce_and_register_asset_event(
+        dag_id="asset-event-producer-decoy",
+        asset=asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+    fired_dag_run_id = 
session.scalar(select(DagRun.id).order_by(DagRun.id.desc()).limit(1))
+    assert fired_dag_run_id is not None
+    decoy.created_dag_run_id = fired_dag_run_id
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_stale_dag_decoy(dag_maker: DagMaker, 
session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    Mirrors :func:`test_partition_cap_reporting_excludes_already_fired_decoy` 
for the fetch
+    query's other predicate: ``DagModel.is_stale.is_(False)``. This test adds 
a decoy pending
+    APDR whose target Dag has since gone stale (removed from its Dag file) — a 
count query
+    missing that filter would wrongly include it, over-counting the backlog 
even though the
+    fetch query (unaffected by this specific drift) never selects it.
+    """
+    cap = 3
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-stale-decoy",
+        asset=Asset(name="asset-cap-stale-decoy"),
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    stale_consumer_dag_id = "cap-consumer-stale-decoy-stale"
+    stale_asset = Asset(name="asset-cap-stale-decoy-stale")
+    with dag_maker(
+        dag_id=stale_consumer_dag_id,
+        schedule=PartitionedAssetTimetable(assets=stale_asset, 
default_partition_mapper=IdentityMapper()),
+        session=session,
+    ):
+        EmptyOperator(task_id="hi")
+    session.commit()
+
+    _produce_and_register_asset_event(
+        dag_id="asset-event-producer-stale-decoy",
+        asset=stale_asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    dm = session.get(DagModel, stale_consumer_dag_id)
+    assert dm is not None
+    dm.is_stale = True
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_lists_all_affected_dag_ids(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The log/audit ``dag_ids`` list surfaces every distinct dag_id in the 
backlog, not just this
+    tick's oldest-cap slice.
+
+    Seven consumer Dags each get one satisfied APDR; with the cap at 6, 
`pending_apdrs` (this
+    tick's oldest six) spans only six distinct dag_ids, but the seventh Dag's 
APDR is still part
+    of the backlog `count()` sees — its dag_id must appear in both the log 
event and the audit
+    row too, even though its Dag run won't be created until the next tick.
+
+    Built inline rather than via :func:`_make_n_satisfied_apdrs` because that 
helper's producer
+    dag_id numbering restarts at 1 on every call, colliding across these seven 
consumer Dags.
+    """
+    apdrs = []
+    for i in range(1, 8):
+        asset = Asset(name=f"asset-cap-many-{i}")
+        with dag_maker(
+            dag_id=f"cap-consumer-many-{i}",
+            schedule=PartitionedAssetTimetable(assets=asset, 
default_partition_mapper=IdentityMapper()),
+            session=session,
+        ):
+            EmptyOperator(task_id="hi")
+        session.commit()
+        apdrs.append(
+            _produce_and_register_asset_event(
+                dag_id=f"asset-event-producer-many-{i}",
+                asset=asset,
+                partition_key=f"k{i}",
+                session=session,
+                dag_maker=dag_maker,
+            )
+        )
+    base = timezone.utcnow()
+    for i, apdr in enumerate(apdrs):
+        apdr.created_at = base + timedelta(seconds=i)
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 6
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    # The dag_ids list is sorted by (-count, dag_id); all seven Dags tie on 
count (one pending
+    # APDR each), so the tie-break is dag_id ascending — a deterministic, 
reproducible order,
+    # not just membership. All seven Dags are expected, including the seventh 
whose APDR is
+    # deferred past this tick's cap.
+    expected_dag_ids = sorted(f"cap-consumer-many-{i}" for i in range(1, 8))
+    assert {
+        "event": "Reached the per-tick cap on pending partitioned Dag runs; 
the remaining backlog "
+        "will be evaluated over subsequent scheduler ticks",
+        "cap": 6,
+        "backlog_total": 7,
+        "dag_ids": expected_dag_ids,
+    } in caplog
+
+    audit_row = session.scalar(select(Log).where(Log.event == "partition Dag 
run cap reached"))
+    assert audit_row is not None
+    assert audit_row.extra is not None
+    note_match = re.search(r"Affected dag_ids \(highest pending count first\): 
(.+)\.$", audit_row.extra)
+    assert note_match is not None
+    assert note_match.group(1).split(", ") == expected_dag_ids
+    assert "and 1 more" not in audit_row.extra
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_truncates_dag_id_list_beyond_max(
+    dag_maker: DagMaker, session: Session, caplog
+):
+    """
+    The log/audit ``dag_ids`` list is capped at 
``MAX_PARTITION_CAP_BACKLOG_DAG_IDS_LOGGED``,
+    ordered by descending per-dag pending count with a dag_id tie-break — not 
alphabetically
+    truncated — and the audit row's prose notes how many dag_ids were dropped 
by the cap.
+
+    Built via :func:`_make_n_consumer_dags_each_with_one_pending_apdr` rather 
than
+    :func:`_make_n_satisfied_apdrs` because that helper's producer dag_id 
numbering restarts at 1
+    on every call, colliding across these 21 consumer Dags — same issue 
documented on
+    :func:`test_partition_cap_reporting_lists_all_affected_dag_ids`.
+    """
+    _make_n_consumer_dags_each_with_one_pending_apdr(
+        name_prefix="trunc-low", n=20, session=session, dag_maker=dag_maker
+    )
+
+    high_asset = Asset(name="asset-cap-trunc-high")
+    with dag_maker(
+        dag_id="cap-consumer-trunc-high",
+        schedule=PartitionedAssetTimetable(assets=high_asset, 
default_partition_mapper=IdentityMapper()),
+        session=session,
+    ):
+        EmptyOperator(task_id="hi")
+    session.commit()
+    for key in ["high-1", "high-2", "high-3"]:
+        _produce_and_register_asset_event(
+            dag_id=f"asset-event-producer-trunc-high-{key}",
+            asset=high_asset,
+            partition_key=key,
+            session=session,
+            dag_maker=dag_maker,
+        )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 5
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    expected_dag_ids = ["cap-consumer-trunc-high"] + 
[f"cap-consumer-trunc-low-{i:02d}" for i in range(1, 20)]
+    assert {
+        "event": "Reached the per-tick cap on pending partitioned Dag runs; 
the remaining backlog "
+        "will be evaluated over subsequent scheduler ticks",
+        "cap": 5,
+        "backlog_total": 23,
+        "dag_ids": expected_dag_ids,
+    } in caplog
+
+    audit_row = session.scalar(select(Log).where(Log.event == "partition Dag 
run cap reached"))
+    assert audit_row is not None
+    assert audit_row.extra is not None
+    assert "and 1 more" in audit_row.extra
+    for dag_id in expected_dag_ids:
+        assert dag_id in audit_row.extra
+    assert "cap-consumer-trunc-low-20" not in audit_row.extra
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_at_exactly_max_dag_ids_no_truncation(
+    dag_maker: DagMaker, session: Session, caplog
+):
+    """
+    Cap-boundary pair to 
:func:`test_partition_cap_reporting_truncates_dag_id_list_beyond_max`:
+    when the backlog's distinct dag_id count equals 
``MAX_PARTITION_CAP_BACKLOG_DAG_IDS_LOGGED``
+    exactly, no dag_id is dropped and the "and N more" suffix is absent.
+
+    Built via :func:`_make_n_consumer_dags_each_with_one_pending_apdr` for the 
same producer
+    dag_id collision reason documented on the neighboring cap tests above.
+    """
+    _make_n_consumer_dags_each_with_one_pending_apdr(
+        name_prefix="atmax", n=20, session=session, dag_maker=dag_maker
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 5
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    expected_dag_ids = sorted(f"cap-consumer-atmax-{i:02d}" for i in range(1, 
21))
+    assert {
+        "event": "Reached the per-tick cap on pending partitioned Dag runs; 
the remaining backlog "
+        "will be evaluated over subsequent scheduler ticks",
+        "cap": 5,
+        "backlog_total": 20,
+        "dag_ids": expected_dag_ids,
+    } in caplog
+
+    audit_row = session.scalar(select(Log).where(Log.event == "partition Dag 
run cap reached"))
+    assert audit_row is not None
+    assert audit_row.extra is not None
+    assert "more" not in audit_row.extra
+    for dag_id in expected_dag_ids:
+        assert dag_id in audit_row.extra
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_backlog_audit_row_written_once_per_episode(
+    dag_maker: DagMaker, session: Session, caplog
+):
+    """
+    The cap-reached audit row is written once per backlog episode, not once 
per tick.
+
+    A persistent backlog re-hits the cap on every tick; only the first tick of 
the episode
+    should write the audit row. Draining below the cap and then hitting it 
again starts a new
+    episode and writes a second row.
+    """
+    asset = Asset(name="asset-cap-episode")
+    apdrs = _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-episode",
+        asset=asset,
+        partition_keys=["k1", "k2", "k3", "k4", "k5"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+    base = timezone.utcnow()
+    for i, apdr in enumerate(apdrs):
+        apdr.created_at = base + timedelta(seconds=i)
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 2
+    caplog.set_level("DEBUG", 
logger="airflow.jobs.scheduler_job_runner.SchedulerJobRunner")
+
+    def _count_audit_log() -> int:
+        return session.scalar(select(func.count()).where(Log.event == 
"partition Dag run cap reached")) or 0
+
+    def _cap_reached_levels() -> list[str]:
+        return [
+            e["log_level"]
+            for e in caplog
+            if e.get("event") == "Reached the per-tick cap on pending 
partitioned Dag runs; the remaining "
+            "backlog will be evaluated over subsequent scheduler ticks"
+        ]
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
1: 5 pending, cap 2
+    assert _cap_reached_levels() == ["warning"]
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
2: 3 pending, still over cap
+    assert _count_audit_log() == 1, "backlog persisted across two ticks; audit 
row must only be written once"
+    assert _cap_reached_levels() == ["warning", "debug"]
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
3: 1 pending, drains below cap
+    assert _count_audit_log() == 1, "draining below the cap must not write 
another audit row"
+    assert _cap_reached_levels() == ["warning", "debug"]
+
+    # A fresh backlog episode for the same consumer: three more satisfied 
APDRs push
+    # pending count back above the cap.
+    new_apdrs = [
+        _produce_and_register_asset_event(
+            dag_id=f"asset-event-producer-episode-{i}",
+            asset=asset,
+            partition_key=key,
+            session=session,
+            dag_maker=dag_maker,
+        )
+        for i, key in enumerate(["k6", "k7", "k8"], start=1)
+    ]
+    base2 = timezone.utcnow()
+    for i, apdr in enumerate(new_apdrs):
+        apdr.created_at = base2 + timedelta(seconds=i)
+    session.commit()
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
4: new episode, 3 pending
+    assert _count_audit_log() == 2, "a fresh backlog episode after draining 
should write a second audit row"
+    assert _cap_reached_levels() == ["warning", "debug", "warning"]
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_backlog_reset_on_full_drain_to_empty(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    Regression test: the backlog-reported flag must reset even when the 
backlog drains to
+    literally zero pending APDRs, not just when it drops to some 
smaller-than-cap-but-still
+    -nonzero count.
+
+    ``_partition_cap_backlog_reported`` used to be cleared only in the 
``len(pending_apdrs) <
+    cap`` branch; a tick whose query returns *zero* rows short-circuits via an 
early ``return
+    set()`` before that branch ever runs. If the last pending APDR of an 
episode gets resolved by
+    something other than this function's own firing (e.g. a concurrent HA 
scheduler winning the
+    ``skip_locked`` race), the very next tick sees an empty ``pending_apdrs`` 
and, without the
+    fix, the stale ``True`` flag survives — so a brand new backlog that 
re-crosses the cap is
+    wrongly logged at ``debug`` and skips writing a second audit row.
+    """
+    asset = Asset(name="asset-cap-drain-empty")
+    apdrs = _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-drain-empty",
+        asset=asset,
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+    base = timezone.utcnow()
+    for i, apdr in enumerate(apdrs):
+        apdr.created_at = base + timedelta(seconds=i)
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 2
+    caplog.set_level("DEBUG", 
logger="airflow.jobs.scheduler_job_runner.SchedulerJobRunner")
+
+    def _count_audit_log() -> int:
+        return session.scalar(select(func.count()).where(Log.event == 
"partition Dag run cap reached")) or 0
+
+    def _cap_reached_levels() -> list[str]:
+        return [
+            e["log_level"]
+            for e in caplog
+            if e.get("event") == "Reached the per-tick cap on pending 
partitioned Dag runs; the remaining "
+            "backlog will be evaluated over subsequent scheduler ticks"
+        ]
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
1: 3 pending, cap 2 -> over cap
+    assert _cap_reached_levels() == ["warning"]
+    assert _count_audit_log() == 1
+    assert runner._partition_cap_backlog_reported is True
+
+    session.refresh(apdrs[0])
+    session.refresh(apdrs[1])
+    session.refresh(apdrs[2])
+    assert apdrs[0].created_dag_run_id is not None
+    assert apdrs[1].created_dag_run_id is not None
+    assert apdrs[2].created_dag_run_id is None
+    # Simulate the one leftover APDR being resolved by something other than 
this function's own
+    # firing this tick (e.g. a concurrent HA scheduler winning the 
`skip_locked` race), so the
+    # *next* tick's query returns zero rows and hits the early-return path 
directly instead of
+    # the `len(pending_apdrs) < cap` branch.
+    apdrs[2].created_dag_run_id = apdrs[0].created_dag_run_id
+    session.commit()
+
+    result = 
runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 2: 0 
pending
+    assert result == set()
+    assert _count_audit_log() == 1, "the zero-pending tick itself must not 
write another audit row"
+    assert runner._partition_cap_backlog_reported is False, (
+        "the backlog-reported flag must reset even on the zero-pending 
early-return path"
+    )
+
+    # A fresh backlog episode for the same consumer Dag: three more satisfied 
APDRs push the
+    # pending count back above the cap.
+    new_apdrs = [
+        _produce_and_register_asset_event(
+            dag_id=f"asset-event-producer-drain-empty-{i}",
+            asset=asset,
+            partition_key=key,
+            session=session,
+            dag_maker=dag_maker,
+        )
+        for i, key in enumerate(["k4", "k5", "k6"], start=1)
+    ]
+    base2 = timezone.utcnow()
+    for i, apdr in enumerate(new_apdrs):
+        apdr.created_at = base2 + timedelta(seconds=i)
+    session.commit()
+
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)  # tick 
3: new episode, 3 pending
+    assert _cap_reached_levels() == ["warning", "warning"], (
+        "a new episode after a full drain-to-zero must warn, not silently log 
at debug"
+    )
+    assert _count_audit_log() == 2, "a new episode after a full drain-to-zero 
must write a second audit row"
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_audit_row_survives_outer_rollback(dag_maker: DagMaker, 
session: Session):
+    """
+    The audit `Log` row is committed in its own session, independent of the 
caller's.
+
+    Rolling back *session* afterwards — as `_create_dagruns_for_dags`'s 
`@retry_db_transaction`
+    would on a `DBAPIError` — must not undo the audit row, and the flag must 
stay ``True``: it
+    reflects a row that is already durably persisted, not in-flight work tied 
to *session*.
+    """
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-rollback",
+        asset=Asset(name="asset-cap-rollback"),
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 2
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+    assert runner._partition_cap_backlog_reported is True
+
+    session.rollback()
+
+    assert runner._partition_cap_backlog_reported is True
+    audit_events = session.scalars(
+        select(Log.event).where(Log.event == "partition Dag run cap reached")
+    ).all()
+    assert audit_events == ["partition Dag run cap reached"]
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
[email protected](
+    "audit_write_error",
+    [
+        pytest.param(IntegrityError("INSERT INTO log", {}, 
Exception("duplicate key")), id="dbapi_error"),
+        pytest.param(
+            # Not a DBAPIError: exercises the broader SQLAlchemyError catch, 
distinct from a
+            # DBAPI-layer failure.
+            InvalidRequestError("session already flushing"),
+            id="non_dbapi_sqlalchemy_error",
+        ),
+    ],
+)
+def test_partition_cap_audit_row_write_failure_is_swallowed(
+    dag_maker: DagMaker, session: Session, caplog, audit_write_error: Exception
+):
+    """
+    A `SQLAlchemyError` while writing the cap-reached audit row must not 
escape the tick.
+
+    The audit write is purely observational, so any SQLAlchemy-layer failure — 
not just a
+    `DBAPIError` — must be caught and logged, leaving 
`_partition_cap_backlog_reported`
+    `False` so the next tick retries the write, and the cap-reached log stays 
at `warning`
+    instead of being downgraded.
+    """
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-audit-failure",
+        asset=Asset(name="asset-cap-audit-failure"),
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    # cap=1 so the 3-item backlog still exceeds the cap on the second tick too 
(3 -> 2 -> 1
+    # remaining), keeping both ticks in the cap-reached branch.
+    runner._max_partition_dag_runs_per_loop = 1
+
+    failing_session_cm = MagicMock()

Review Comment:
   Can this be `MagicMock(spec=Session)` on the `__enter__` return value (or 
`create_autospec`)? Also, the point of swallowing the error is that the tick's 
DagRun creation still goes ahead, but nothing asserts that. A broken `except: 
return set()` would pass this test too. Asserting that two of the three APDRs 
got a `created_dag_run_id` would cover it.



##########
airflow-core/src/airflow/jobs/scheduler_job_runner.py:
##########
@@ -2349,8 +2363,95 @@ def _create_dagruns_for_partitioned_asset_dags(self, 
session: Session) -> set[st
             )
         ).all()
         if not pending_apdrs:
+            # Backlog fully drained: the next tick that re-crosses the cap is 
a new episode,
+            # not a continuation of whatever was reported before.
+            self._partition_cap_backlog_reported = False
             return set()
 
+        sorted_dag_ids: list[str] = []
+        if len(pending_apdrs) >= self._max_partition_dag_runs_per_loop:
+            # A full fetch alone can't tell us whether that's the entire 
backlog or just
+            # this tick's slice of a larger one, so we only pay for this query 
then.
+            # Per-Dag counts across the *whole* backlog, not just this tick's 
oldest-cap
+            # slice (`pending_apdrs`) — a Dag whose partitions haven't reached 
the front of
+            # the FIFO queue yet would otherwise be missing from the log/audit 
row until its
+            # turn comes up.
+            backlogs_per_dag: dict[str, int] = {
+                dag_id: count
+                for dag_id, count in session.execute(
+                    select(AssetPartitionDagRun.target_dag_id, func.count())
+                    .select_from(AssetPartitionDagRun)
+                    .join(DagModel, DagModel.dag_id == 
AssetPartitionDagRun.target_dag_id)
+                    .where(
+                        AssetPartitionDagRun.created_dag_run_id.is_(None),
+                        DagModel.is_stale.is_(False),
+                    )
+                    .group_by(AssetPartitionDagRun.target_dag_id)
+                )
+            }
+            backlog_total = sum(backlogs_per_dag.values())
+            # No SQL-side ORDER BY: it would have no dag_id tie-break, and 
relying on dict
+            # insertion order mirroring DB row order across backends (SQLite 
in tests vs
+            # Postgres/MySQL in production) isn't a contract worth trusting — 
sort explicitly.
+            sorted_dag_ids = sorted(backlogs_per_dag, key=lambda dag_id: 
(-backlogs_per_dag[dag_id], dag_id))
+        else:
+            backlog_total = len(pending_apdrs)
+            self._partition_cap_backlog_reported = False

Review Comment:
   With `skip_locked=True`, a short fetch here doesn't mean the backlog 
drained. In HA with a backlog of 600, scheduler A locks 500 and B gets 100, so 
B clears its flag. On the next tick B wins the full slice and writes a fresh 
warning and audit row for the same episode, and this can repeat for as long as 
the backlog sits near the cap. So the comment at L385 ("at most one audit row 
per episode") doesn't hold. Could the reset happen only on an empty fetch, or 
after an unlocked count shows total <= cap? That would also make this reset and 
the one at L2453 collapse into one.



##########
airflow-core/src/airflow/jobs/scheduler_job_runner.py:
##########
@@ -2349,8 +2363,95 @@ def _create_dagruns_for_partitioned_asset_dags(self, 
session: Session) -> set[st
             )
         ).all()
         if not pending_apdrs:
+            # Backlog fully drained: the next tick that re-crosses the cap is 
a new episode,
+            # not a continuation of whatever was reported before.
+            self._partition_cap_backlog_reported = False
             return set()
 
+        sorted_dag_ids: list[str] = []
+        if len(pending_apdrs) >= self._max_partition_dag_runs_per_loop:
+            # A full fetch alone can't tell us whether that's the entire 
backlog or just
+            # this tick's slice of a larger one, so we only pay for this query 
then.
+            # Per-Dag counts across the *whole* backlog, not just this tick's 
oldest-cap
+            # slice (`pending_apdrs`) — a Dag whose partitions haven't reached 
the front of
+            # the FIFO queue yet would otherwise be missing from the log/audit 
row until its
+            # turn comes up.
+            backlogs_per_dag: dict[str, int] = {
+                dag_id: count
+                for dag_id, count in session.execute(
+                    select(AssetPartitionDagRun.target_dag_id, func.count())
+                    .select_from(AssetPartitionDagRun)
+                    .join(DagModel, DagModel.dag_id == 
AssetPartitionDagRun.target_dag_id)
+                    .where(
+                        AssetPartitionDagRun.created_dag_run_id.is_(None),
+                        DagModel.is_stale.is_(False),

Review Comment:
   This count is missing the `DagModel.is_paused.is_(False)` and 
`DagModel.is_draining.is_(False)` filters that the fetch above has. Say a 
paused Dag has 2000 pending APDRs and exactly 500 other APDRs are eligible: the 
fetch returns 500, this returns 2500, and we warn about a backlog that doesn't 
exist, with the paused Dag listed first in `dag_ids` and in the audit row. 
Could the two queries share one WHERE clause so they can't drift? The decoy 
tests cover fired and stale, so adding paused and draining cases (parametrizing 
the two existing decoy tests) would lock it in.



##########
airflow-core/tests/unit/jobs/test_scheduler_job.py:
##########
@@ -13244,6 +13296,594 @@ def 
test_partition_cap_at_n_minus_one_leaves_one_pending(dag_maker: DagMaker, se
     assert partition_dags == {"cap-consumer-n-minus-one"}
 
 
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
[email protected](
+    ("cap", "partition_keys", "expect_cap_hit"),
+    [
+        pytest.param(2, ["k1", "k2", "k3"], True, id="over-cap"),
+        pytest.param(3, ["k1", "k2"], False, id="under-cap"),
+        # Regression guard for the ``backlog_total > cap`` boundary: 
``backlog_total == cap``
+        # must not be misreported as a backlog (that would be the "exactly 
cap, nothing left"
+        # case wrongly reported as "more than cap, backlog remains").
+        pytest.param(3, ["k1", "k2", "k3"], False, id="exactly-cap"),
+    ],
+)
+def test_partition_cap_reporting(
+    dag_maker: DagMaker,
+    session: Session,
+    caplog,
+    cap: int,
+    partition_keys: list[str],
+    expect_cap_hit: bool,
+):
+    """Only a pending count strictly above the cap reports a backlog, via log 
and audit `Log` row."""
+    suffix = f"{cap}-{len(partition_keys)}"
+    _make_n_satisfied_apdrs(
+        consumer_dag_id=f"cap-consumer-{suffix}",
+        asset=Asset(name=f"asset-cap-{suffix}"),
+        partition_keys=partition_keys,
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    dag_id = f"cap-consumer-{suffix}"
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    if expect_cap_hit:
+        assert {
+            "event": "Reached the per-tick cap on pending partitioned Dag 
runs; the remaining backlog "
+            "will be evaluated over subsequent scheduler ticks",
+            "cap": cap,
+            "backlog_total": 3,
+            "dag_ids": [dag_id],
+        } in caplog
+        assert [row.event for row in audit_rows] == ["partition Dag run cap 
reached"]
+        extra = audit_rows[0].extra
+        assert extra is not None
+        assert f"Affected dag_ids (highest pending count first): {dag_id}." in 
extra
+    else:
+        assert "Reached the per-tick cap on pending partitioned Dag runs" not 
in caplog
+        assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_already_fired_decoy(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    ``exactly-cap`` in :func:`test_partition_cap_reporting` cannot catch a 
filter drift in the
+    count query (e.g. a dropped ``created_dag_run_id.is_(None)`` filter): 
every row in that
+    table at that point is a genuine, unfired, non-stale APDR, so a looser 
filter has nothing
+    extra to over-count. This test adds a decoy APDR that already fired
+    (``created_dag_run_id`` set) — a naive/unfiltered ``count()`` would 
wrongly include it.
+    """
+    cap = 3
+    asset = Asset(name="asset-cap-decoy")
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-decoy",
+        asset=asset,
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+    decoy = _produce_and_register_asset_event(
+        dag_id="asset-event-producer-decoy",
+        asset=asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+    fired_dag_run_id = 
session.scalar(select(DagRun.id).order_by(DagRun.id.desc()).limit(1))
+    assert fired_dag_run_id is not None
+    decoy.created_dag_run_id = fired_dag_run_id
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_stale_dag_decoy(dag_maker: DagMaker, 
session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    Mirrors :func:`test_partition_cap_reporting_excludes_already_fired_decoy` 
for the fetch
+    query's other predicate: ``DagModel.is_stale.is_(False)``. This test adds 
a decoy pending
+    APDR whose target Dag has since gone stale (removed from its Dag file) — a 
count query
+    missing that filter would wrongly include it, over-counting the backlog 
even though the
+    fetch query (unaffected by this specific drift) never selects it.
+    """
+    cap = 3
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-stale-decoy",
+        asset=Asset(name="asset-cap-stale-decoy"),
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    stale_consumer_dag_id = "cap-consumer-stale-decoy-stale"
+    stale_asset = Asset(name="asset-cap-stale-decoy-stale")
+    with dag_maker(
+        dag_id=stale_consumer_dag_id,
+        schedule=PartitionedAssetTimetable(assets=stale_asset, 
default_partition_mapper=IdentityMapper()),
+        session=session,
+    ):
+        EmptyOperator(task_id="hi")
+    session.commit()
+
+    _produce_and_register_asset_event(
+        dag_id="asset-event-producer-stale-decoy",
+        asset=stale_asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    dm = session.get(DagModel, stale_consumer_dag_id)
+    assert dm is not None
+    dm.is_stale = True
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_lists_all_affected_dag_ids(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The log/audit ``dag_ids`` list surfaces every distinct dag_id in the 
backlog, not just this
+    tick's oldest-cap slice.
+
+    Seven consumer Dags each get one satisfied APDR; with the cap at 6, 
`pending_apdrs` (this
+    tick's oldest six) spans only six distinct dag_ids, but the seventh Dag's 
APDR is still part
+    of the backlog `count()` sees — its dag_id must appear in both the log 
event and the audit
+    row too, even though its Dag run won't be created until the next tick.
+
+    Built inline rather than via :func:`_make_n_satisfied_apdrs` because that 
helper's producer
+    dag_id numbering restarts at 1 on every call, colliding across these seven 
consumer Dags.
+    """
+    apdrs = []
+    for i in range(1, 8):

Review Comment:
   This loop is what `_make_n_consumer_dags_each_with_one_pending_apdr` does. 
Could it use the helper (with the `created_at` spacing applied afterwards) and 
drop the docstring paragraph about producer numbering? `_count_audit_log` and 
`_cap_reached_levels` are also copied verbatim between the two episode tests 
and could be module-level helpers.



##########
airflow-core/tests/unit/jobs/test_scheduler_job.py:
##########
@@ -13244,6 +13296,594 @@ def 
test_partition_cap_at_n_minus_one_leaves_one_pending(dag_maker: DagMaker, se
     assert partition_dags == {"cap-consumer-n-minus-one"}
 
 
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
[email protected](
+    ("cap", "partition_keys", "expect_cap_hit"),
+    [
+        pytest.param(2, ["k1", "k2", "k3"], True, id="over-cap"),
+        pytest.param(3, ["k1", "k2"], False, id="under-cap"),
+        # Regression guard for the ``backlog_total > cap`` boundary: 
``backlog_total == cap``
+        # must not be misreported as a backlog (that would be the "exactly 
cap, nothing left"
+        # case wrongly reported as "more than cap, backlog remains").
+        pytest.param(3, ["k1", "k2", "k3"], False, id="exactly-cap"),
+    ],
+)
+def test_partition_cap_reporting(
+    dag_maker: DagMaker,
+    session: Session,
+    caplog,
+    cap: int,
+    partition_keys: list[str],
+    expect_cap_hit: bool,
+):
+    """Only a pending count strictly above the cap reports a backlog, via log 
and audit `Log` row."""
+    suffix = f"{cap}-{len(partition_keys)}"
+    _make_n_satisfied_apdrs(
+        consumer_dag_id=f"cap-consumer-{suffix}",
+        asset=Asset(name=f"asset-cap-{suffix}"),
+        partition_keys=partition_keys,
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    dag_id = f"cap-consumer-{suffix}"
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    if expect_cap_hit:
+        assert {
+            "event": "Reached the per-tick cap on pending partitioned Dag 
runs; the remaining backlog "
+            "will be evaluated over subsequent scheduler ticks",
+            "cap": cap,
+            "backlog_total": 3,
+            "dag_ids": [dag_id],
+        } in caplog
+        assert [row.event for row in audit_rows] == ["partition Dag run cap 
reached"]
+        extra = audit_rows[0].extra
+        assert extra is not None
+        assert f"Affected dag_ids (highest pending count first): {dag_id}." in 
extra
+    else:
+        assert "Reached the per-tick cap on pending partitioned Dag runs" not 
in caplog
+        assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_already_fired_decoy(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    ``exactly-cap`` in :func:`test_partition_cap_reporting` cannot catch a 
filter drift in the
+    count query (e.g. a dropped ``created_dag_run_id.is_(None)`` filter): 
every row in that
+    table at that point is a genuine, unfired, non-stale APDR, so a looser 
filter has nothing
+    extra to over-count. This test adds a decoy APDR that already fired
+    (``created_dag_run_id`` set) — a naive/unfiltered ``count()`` would 
wrongly include it.
+    """
+    cap = 3
+    asset = Asset(name="asset-cap-decoy")
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-decoy",
+        asset=asset,
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+    decoy = _produce_and_register_asset_event(
+        dag_id="asset-event-producer-decoy",
+        asset=asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+    fired_dag_run_id = 
session.scalar(select(DagRun.id).order_by(DagRun.id.desc()).limit(1))
+    assert fired_dag_run_id is not None
+    decoy.created_dag_run_id = fired_dag_run_id
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_excludes_stale_dag_decoy(dag_maker: DagMaker, 
session: Session, caplog):
+    """
+    The ``backlog_total`` count query's ``WHERE`` clause must mirror the main 
query's.
+
+    Mirrors :func:`test_partition_cap_reporting_excludes_already_fired_decoy` 
for the fetch
+    query's other predicate: ``DagModel.is_stale.is_(False)``. This test adds 
a decoy pending
+    APDR whose target Dag has since gone stale (removed from its Dag file) — a 
count query
+    missing that filter would wrongly include it, over-counting the backlog 
even though the
+    fetch query (unaffected by this specific drift) never selects it.
+    """
+    cap = 3
+    _make_n_satisfied_apdrs(
+        consumer_dag_id="cap-consumer-stale-decoy",
+        asset=Asset(name="asset-cap-stale-decoy"),
+        partition_keys=["k1", "k2", "k3"],
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    stale_consumer_dag_id = "cap-consumer-stale-decoy-stale"
+    stale_asset = Asset(name="asset-cap-stale-decoy-stale")
+    with dag_maker(
+        dag_id=stale_consumer_dag_id,
+        schedule=PartitionedAssetTimetable(assets=stale_asset, 
default_partition_mapper=IdentityMapper()),
+        session=session,
+    ):
+        EmptyOperator(task_id="hi")
+    session.commit()
+
+    _produce_and_register_asset_event(
+        dag_id="asset-event-producer-stale-decoy",
+        asset=stale_asset,
+        partition_key="decoy",
+        session=session,
+        dag_maker=dag_maker,
+    )
+
+    dm = session.get(DagModel, stale_consumer_dag_id)
+    assert dm is not None
+    dm.is_stale = True
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = cap
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    assert "Reached the per-tick cap on pending partitioned Dag runs" not in 
caplog
+    audit_rows = session.scalars(select(Log).where(Log.event == "partition Dag 
run cap reached")).all()
+    assert audit_rows == []
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_lists_all_affected_dag_ids(dag_maker: 
DagMaker, session: Session, caplog):
+    """
+    The log/audit ``dag_ids`` list surfaces every distinct dag_id in the 
backlog, not just this
+    tick's oldest-cap slice.
+
+    Seven consumer Dags each get one satisfied APDR; with the cap at 6, 
`pending_apdrs` (this
+    tick's oldest six) spans only six distinct dag_ids, but the seventh Dag's 
APDR is still part
+    of the backlog `count()` sees — its dag_id must appear in both the log 
event and the audit
+    row too, even though its Dag run won't be created until the next tick.
+
+    Built inline rather than via :func:`_make_n_satisfied_apdrs` because that 
helper's producer
+    dag_id numbering restarts at 1 on every call, colliding across these seven 
consumer Dags.
+    """
+    apdrs = []
+    for i in range(1, 8):
+        asset = Asset(name=f"asset-cap-many-{i}")
+        with dag_maker(
+            dag_id=f"cap-consumer-many-{i}",
+            schedule=PartitionedAssetTimetable(assets=asset, 
default_partition_mapper=IdentityMapper()),
+            session=session,
+        ):
+            EmptyOperator(task_id="hi")
+        session.commit()
+        apdrs.append(
+            _produce_and_register_asset_event(
+                dag_id=f"asset-event-producer-many-{i}",
+                asset=asset,
+                partition_key=f"k{i}",
+                session=session,
+                dag_maker=dag_maker,
+            )
+        )
+    base = timezone.utcnow()
+    for i, apdr in enumerate(apdrs):
+        apdr.created_at = base + timedelta(seconds=i)
+    session.commit()
+
+    runner = SchedulerJobRunner(
+        job=Job(job_type=SchedulerJobRunner.job_type), 
executors=[MockExecutor(do_update=False)]
+    )
+    runner._max_partition_dag_runs_per_loop = 6
+    runner._create_dagruns_for_partitioned_asset_dags(session=session)
+
+    # The dag_ids list is sorted by (-count, dag_id); all seven Dags tie on 
count (one pending
+    # APDR each), so the tie-break is dag_id ascending — a deterministic, 
reproducible order,
+    # not just membership. All seven Dags are expected, including the seventh 
whose APDR is
+    # deferred past this tick's cap.
+    expected_dag_ids = sorted(f"cap-consumer-many-{i}" for i in range(1, 8))
+    assert {
+        "event": "Reached the per-tick cap on pending partitioned Dag runs; 
the remaining backlog "
+        "will be evaluated over subsequent scheduler ticks",
+        "cap": 6,
+        "backlog_total": 7,
+        "dag_ids": expected_dag_ids,
+    } in caplog
+
+    audit_row = session.scalar(select(Log).where(Log.event == "partition Dag 
run cap reached"))
+    assert audit_row is not None
+    assert audit_row.extra is not None
+    note_match = re.search(r"Affected dag_ids \(highest pending count first\): 
(.+)\.$", audit_row.extra)
+    assert note_match is not None
+    assert note_match.group(1).split(", ") == expected_dag_ids
+    assert "and 1 more" not in audit_row.extra
+
+
[email protected]_serialized_dag
[email protected]("clear_asset_partition_rows", "clear_audit_log_rows")
+def test_partition_cap_reporting_truncates_dag_id_list_beyond_max(
+    dag_maker: DagMaker, session: Session, caplog
+):
+    """
+    The log/audit ``dag_ids`` list is capped at 
``MAX_PARTITION_CAP_BACKLOG_DAG_IDS_LOGGED``,
+    ordered by descending per-dag pending count with a dag_id tie-break — not 
alphabetically
+    truncated — and the audit row's prose notes how many dag_ids were dropped 
by the cap.
+
+    Built via :func:`_make_n_consumer_dags_each_with_one_pending_apdr` rather 
than
+    :func:`_make_n_satisfied_apdrs` because that helper's producer dag_id 
numbering restarts at 1
+    on every call, colliding across these 21 consumer Dags — same issue 
documented on
+    :func:`test_partition_cap_reporting_lists_all_affected_dag_ids`.
+    """
+    _make_n_consumer_dags_each_with_one_pending_apdr(
+        name_prefix="trunc-low", n=20, session=session, dag_maker=dag_maker
+    )
+
+    high_asset = Asset(name="asset-cap-trunc-high")
+    with dag_maker(
+        dag_id="cap-consumer-trunc-high",

Review Comment:
   `cap-consumer-trunc-high` already sorts alphabetically before 
`cap-consumer-trunc-low-NN`, so this test still passes if the sort is by name 
only and the count key is dropped. Renaming it to something like 
`cap-consumer-trunc-zz-high` would make it actually check the descending-count 
ordering.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to