ashb commented on code in PR #74344:
URL: https://github.com/apache/airflow/pull/74344#discussion_r4208955585


##########
airflow-core/src/airflow/models/dynamic_region.py:
##########
@@ -78,3 +100,118 @@ class DynamicRegion(Base):
         Index("idx_dynamic_region_slot", dag_id, run_id, node_id, 
parent_region_id, parent_region_index),
         Index("idx_dynamic_region_parent_region_id", parent_region_id),
     )
+
+
+def resolve_current_producers(
+    *,
+    dag_id: str,
+    run_id: str,
+    task_id: str,
+    is_mapped: bool,
+    context: ProducerContext | None = None,
+    map_indexes: int | Collection[int] | None = None,
+    region_id: UUID | None = None,
+    region_index: int | None = None,
+    session: Session,
+) -> tuple[TaskInstance, ...]:
+    """Resolve the live producer task instances whose data the caller reads by 
task instance UUID."""
+    from airflow.models.taskinstance import TaskInstance
+
+    if region_index is not None and region_id is None:
+        raise ValueError("region_index requires an explicit producer 
region_id")
+    if context and context.previous_iteration and context.loop_node_id is None:
+        raise ValueError("Previous-iteration lookup requires a loop context")
+    query = select(TaskInstance).where(
+        TaskInstance.dag_id == dag_id,
+        TaskInstance.run_id == run_id,
+        TaskInstance.task_id == task_id,
+        TaskInstance.working_set.is_(True),
+    )
+    if region_id is not None:
+        query = query.where(TaskInstance.region_id == region_id)
+    if region_index is not None:
+        query = query.where(TaskInstance.region_index == region_index)
+    candidates = session.scalars(query).all()

Review Comment:
   I'm going to look at fixing this via ` `WITH RECURSIVE` cte (works across 
all 3 dbs)



-- 
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