kaxil commented on code in PR #74344:
URL: https://github.com/apache/airflow/pull/74344#discussion_r4222223278
##########
airflow-core/src/airflow/models/xcom.py:
##########
@@ -345,17 +353,28 @@ def get_many(
raise ValueError(f"XCom key must be a non-empty string. Received:
{key!r}")
if not run_id:
raise ValueError(f"run_id must be passed. Passed run_id={run_id}")
- statement = build_xcom_read_query(
- producer_ids=select_producers(
+ if producer_ids is not None:
Review Comment:
`try_number` still gets through this guard. With `producer_ids` set it adds
no try filter (that only happens in `select_producers`), but line 383 still
sets `include_all_attempts=True`, so `get_many(producer_ids=..., try_number=2)`
reads whatever attempts the ids name rather than try 2. Adding `try_number` to
the rejected combinations would close it.
##########
airflow-core/src/airflow/models/dagrun.py:
##########
@@ -1735,13 +1739,13 @@ def _expand_mapped_task_if_needed(ti: TI) ->
Iterable[TI] | None:
# Check dependencies.
expansion_happened = False
# Set of task ids for which was already done
_revise_map_indexes_if_mapped
- revised_map_index_task_ids: set[str] = set()
+ revised_map_index_task_ids: set[tuple[str, UUID]] = set()
Review Comment:
The key is `(task_id, region_id)` now, but the comment above still says "Set
of task ids", the one at line 1767 says "once per task id", and the name is
still `revised_map_index_task_ids`. Worth updating them to say per task and
region?
##########
airflow-core/tests/unit/models/test_dynamic_region.py:
##########
@@ -0,0 +1,449 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from uuid import uuid4
+
+import pytest
+from sqlalchemy import event, select
+
+from airflow._shared.timezones import timezone
+from airflow.models.dynamic_region import (
+ AmbiguousProducerError,
+ DynamicRegion,
+ ProducerContext,
+ resolve_current_producers,
+)
+from airflow.models.taskinstance import TaskInstance
+from airflow.models.xcom import XComModel
+from airflow.providers.standard.operators.empty import EmptyOperator
+from airflow.providers.standard.operators.python import PythonOperator
+from airflow.utils.state import TaskInstanceState
+
+from tests_common.test_utils.db import clear_db_runs
+
+pytestmark = pytest.mark.db_test
+
+
[email protected](autouse=True)
+def clean_db():
+ clear_db_runs()
+ yield
+ clear_db_runs()
+
+
[email protected]
+def regional_tis(dag_maker, session):
+ with dag_maker(serialized=True):
+ task = EmptyOperator(task_id="task")
+ dr = dag_maker.create_dagrun()
+ original = dr.task_instances[0]
+ other = TaskInstance(task=task, run_id=dr.run_id,
dag_version_id=original.dag_version_id)
+ other.region_id = uuid4()
+ session.add(other)
+ session.flush()
+ return original, other
+
+
[email protected]
+def producer_tis(dag_maker, session):
+ with dag_maker(serialized=True) as dag:
+ for task_id in ("producer", "consumer", "outside", "mapped"):
+ EmptyOperator(task_id=task_id)
+ dr = dag_maker.create_dagrun()
+ regions = []
+ for _ in range(3):
+ region = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id,
node_id="loop")
+ if regions:
+ region.forked_from_region_id = regions[-1].id
+ session.add(region)
+ session.flush()
+ regions.append(region)
+ tis = {ti.task_id: ti for ti in dr.task_instances}
+ tis["producer"].region_id = regions[0].id
+ tis["producer"].region_index = 2
+ tis["consumer"].region_id = regions[2].id
+ tis["consumer"].region_index = 2
+ previous = TaskInstance(
+ dag.get_task("producer"),
+ tis["producer"].dag_version_id,
+ run_id=dr.run_id,
+ map_index=1,
+ region_id=regions[0].id,
+ )
+ session.add(previous)
+ session.flush()
+ return tis, regions, previous
+
+
[email protected]("previous_iteration", [False, True])
+def test_resolve_retained_producer_across_repeated_forks(producer_tis,
session, previous_iteration):
+ tis, regions, previous = producer_tis
+ consumer = tis["consumer"]
+ selected = resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop", previous_iteration),
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [previous.id if previous_iteration
else tis["producer"].id]
+
+
+def test_resolve_producer_outside_the_loop(producer_tis, session):
+ tis, _, _ = producer_tis
+ consumer = tis["consumer"]
+ selected = resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="outside",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index),
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [tis["outside"].id]
+
+
+def test_resolver_never_revives_archived_producer(producer_tis, session):
+ tis, _, _ = producer_tis
+ producer, consumer = tis["producer"], tis["consumer"]
+ producer.archive(reason="test", session=session)
+ assert (
+ resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop"),
+ session=session,
+ )
+ == ()
+ )
+
+
+def test_resolver_rejects_ambiguous_live_producers(producer_tis, session):
+ tis, regions, _ = producer_tis
+ producer, consumer = tis["producer"], tis["consumer"]
+ other = TaskInstance(
+ producer.task,
+ producer.dag_version_id,
+ run_id=producer.run_id,
+ map_index=producer.map_index,
+ region_id=regions[1].id,
+ )
+ session.add(other)
+ session.flush()
+ with pytest.raises(AmbiguousProducerError):
+ resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop"),
+ session=session,
+ )
+
+
[email protected]("mapped_caller", [False, True])
+def test_mapped_producer_scope_precedes_index_selection(producer_tis, session,
mapped_caller):
+ tis, regions, _ = producer_tis
+ consumer, mapped = tis["consumer"], tis["mapped"]
+ children = []
+ for parent, iteration in ((regions[0], 2), (regions[2], 2), (regions[2],
3)):
+ child = DynamicRegion(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ node_id="mapped",
+ parent_region_id=parent.id,
+ parent_region_index=iteration,
+ )
+ session.add(child)
+ session.flush()
+ children.append(child)
+ mapped.region_id, mapped.region_index = children[0].id, 0
+ second = TaskInstance(
+ mapped.task, mapped.dag_version_id, run_id=mapped.run_id, map_index=1,
region_id=children[1].id
+ )
+ wrong_iteration = TaskInstance(
+ mapped.task, mapped.dag_version_id, run_id=mapped.run_id, map_index=0,
region_id=children[2].id
+ )
+ session.add_all([second, wrong_iteration])
+ if mapped_caller:
+ caller_region = DynamicRegion(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ node_id="consumer",
+ parent_region_id=regions[2].id,
+ parent_region_index=2,
+ )
+ session.add(caller_region)
+ session.flush()
+ consumer.region_id, consumer.region_index = caller_region.id, 5
+ session.flush()
+ context = ProducerContext(consumer.region_id, consumer.region_index,
"loop")
+ selected = resolve_current_producers(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ task_id="mapped",
+ is_mapped=True,
+ context=context,
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [mapped.id, second.id]
+ selected = resolve_current_producers(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ task_id="mapped",
+ is_mapped=True,
+ context=context,
+ map_indexes=1,
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [second.id]
+
+
[email protected]("depth", [1, 2, 3])
+def
test_loop_iteration_lookup_loads_only_the_selected_iterations_rows(producer_tis,
session, depth):
Review Comment:
Nothing here pins the `parent_region_id.in_(select(family.c.id))` condition
in the nested CTE anchor (`dynamic_region.py:196`). Every loop-position test in
this file and in `test_xcom_arg.py` builds a single `node_id="loop"` family, so
`parent_region_index == iteration` alone already picks the right rows, and
dropping the family condition would leave these tests green. That condition is
what stops iteration k of one loop pulling in another family's iteration k.
Could you add a case with a second, unrelated loop root that has a nested
mapped region at the same `parent_region_index`, and assert its TI is neither
selected nor loaded?
##########
airflow-core/tests/unit/models/test_dynamic_region.py:
##########
@@ -0,0 +1,449 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+from __future__ import annotations
+
+from uuid import uuid4
+
+import pytest
+from sqlalchemy import event, select
+
+from airflow._shared.timezones import timezone
+from airflow.models.dynamic_region import (
+ AmbiguousProducerError,
+ DynamicRegion,
+ ProducerContext,
+ resolve_current_producers,
+)
+from airflow.models.taskinstance import TaskInstance
+from airflow.models.xcom import XComModel
+from airflow.providers.standard.operators.empty import EmptyOperator
+from airflow.providers.standard.operators.python import PythonOperator
+from airflow.utils.state import TaskInstanceState
+
+from tests_common.test_utils.db import clear_db_runs
+
+pytestmark = pytest.mark.db_test
+
+
[email protected](autouse=True)
+def clean_db():
+ clear_db_runs()
+ yield
+ clear_db_runs()
+
+
[email protected]
+def regional_tis(dag_maker, session):
+ with dag_maker(serialized=True):
+ task = EmptyOperator(task_id="task")
+ dr = dag_maker.create_dagrun()
+ original = dr.task_instances[0]
+ other = TaskInstance(task=task, run_id=dr.run_id,
dag_version_id=original.dag_version_id)
+ other.region_id = uuid4()
+ session.add(other)
+ session.flush()
+ return original, other
+
+
[email protected]
+def producer_tis(dag_maker, session):
+ with dag_maker(serialized=True) as dag:
+ for task_id in ("producer", "consumer", "outside", "mapped"):
+ EmptyOperator(task_id=task_id)
+ dr = dag_maker.create_dagrun()
+ regions = []
+ for _ in range(3):
+ region = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id,
node_id="loop")
+ if regions:
+ region.forked_from_region_id = regions[-1].id
+ session.add(region)
+ session.flush()
+ regions.append(region)
+ tis = {ti.task_id: ti for ti in dr.task_instances}
+ tis["producer"].region_id = regions[0].id
+ tis["producer"].region_index = 2
+ tis["consumer"].region_id = regions[2].id
+ tis["consumer"].region_index = 2
+ previous = TaskInstance(
+ dag.get_task("producer"),
+ tis["producer"].dag_version_id,
+ run_id=dr.run_id,
+ map_index=1,
+ region_id=regions[0].id,
+ )
+ session.add(previous)
+ session.flush()
+ return tis, regions, previous
+
+
[email protected]("previous_iteration", [False, True])
+def test_resolve_retained_producer_across_repeated_forks(producer_tis,
session, previous_iteration):
+ tis, regions, previous = producer_tis
+ consumer = tis["consumer"]
+ selected = resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop", previous_iteration),
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [previous.id if previous_iteration
else tis["producer"].id]
+
+
+def test_resolve_producer_outside_the_loop(producer_tis, session):
+ tis, _, _ = producer_tis
+ consumer = tis["consumer"]
+ selected = resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="outside",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index),
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [tis["outside"].id]
+
+
+def test_resolver_never_revives_archived_producer(producer_tis, session):
+ tis, _, _ = producer_tis
+ producer, consumer = tis["producer"], tis["consumer"]
+ producer.archive(reason="test", session=session)
+ assert (
+ resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop"),
+ session=session,
+ )
+ == ()
+ )
+
+
+def test_resolver_rejects_ambiguous_live_producers(producer_tis, session):
+ tis, regions, _ = producer_tis
+ producer, consumer = tis["producer"], tis["consumer"]
+ other = TaskInstance(
+ producer.task,
+ producer.dag_version_id,
+ run_id=producer.run_id,
+ map_index=producer.map_index,
+ region_id=regions[1].id,
+ )
+ session.add(other)
+ session.flush()
+ with pytest.raises(AmbiguousProducerError):
+ resolve_current_producers(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ task_id="producer",
+ is_mapped=False,
+ context=ProducerContext(consumer.region_id, consumer.region_index,
"loop"),
+ session=session,
+ )
+
+
[email protected]("mapped_caller", [False, True])
+def test_mapped_producer_scope_precedes_index_selection(producer_tis, session,
mapped_caller):
+ tis, regions, _ = producer_tis
+ consumer, mapped = tis["consumer"], tis["mapped"]
+ children = []
+ for parent, iteration in ((regions[0], 2), (regions[2], 2), (regions[2],
3)):
+ child = DynamicRegion(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ node_id="mapped",
+ parent_region_id=parent.id,
+ parent_region_index=iteration,
+ )
+ session.add(child)
+ session.flush()
+ children.append(child)
+ mapped.region_id, mapped.region_index = children[0].id, 0
+ second = TaskInstance(
+ mapped.task, mapped.dag_version_id, run_id=mapped.run_id, map_index=1,
region_id=children[1].id
+ )
+ wrong_iteration = TaskInstance(
+ mapped.task, mapped.dag_version_id, run_id=mapped.run_id, map_index=0,
region_id=children[2].id
+ )
+ session.add_all([second, wrong_iteration])
+ if mapped_caller:
+ caller_region = DynamicRegion(
+ dag_id=consumer.dag_id,
+ run_id=consumer.run_id,
+ node_id="consumer",
+ parent_region_id=regions[2].id,
+ parent_region_index=2,
+ )
+ session.add(caller_region)
+ session.flush()
+ consumer.region_id, consumer.region_index = caller_region.id, 5
+ session.flush()
+ context = ProducerContext(consumer.region_id, consumer.region_index,
"loop")
+ selected = resolve_current_producers(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ task_id="mapped",
+ is_mapped=True,
+ context=context,
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [mapped.id, second.id]
+ selected = resolve_current_producers(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ task_id="mapped",
+ is_mapped=True,
+ context=context,
+ map_indexes=1,
+ session=session,
+ )
+ assert [ti.id for ti in selected] == [second.id]
+
+
[email protected]("depth", [1, 2, 3])
+def
test_loop_iteration_lookup_loads_only_the_selected_iterations_rows(producer_tis,
session, depth):
+ tis, regions, _ = producer_tis
+ consumer, mapped = tis["consumer"], tis["mapped"]
+ expected = []
+ for iteration in range(4):
+ parent, parent_index = regions[iteration % 3], iteration
+ for level in range(depth):
+ nested = DynamicRegion(
+ dag_id=mapped.dag_id,
+ run_id=mapped.run_id,
+ node_id=f"mapped{level}",
+ parent_region_id=parent.id,
+ parent_region_index=parent_index,
+ )
+ session.add(nested)
+ session.flush()
+ parent, parent_index = nested, 0
+ for map_index in range(3):
+ ti = TaskInstance(
+ mapped.task,
+ mapped.dag_version_id,
+ run_id=mapped.run_id,
+ map_index=map_index,
+ region_id=parent.id,
+ )
+ session.add(ti)
+ if iteration == consumer.region_index:
+ expected.append(ti)
+ session.flush()
+ expected_ids = {ti.id for ti in expected}
+ session.expunge_all()
+ loaded = []
+ event.listen(session, "loaded_as_persistent", lambda _, instance:
loaded.append(instance))
Review Comment:
This listener is never removed. The `session` fixture comes from
`create_session()`, which hands back the thread-local scoped session and only
`close()`s it, so the listener stays attached for every later test on that
worker. A `try`/`finally` with `event.remove(...)` would keep it local to this
test.
##########
airflow-core/tests/unit/models/test_xcom_arg.py:
##########
@@ -266,3 +276,190 @@ def consume(value): ...
session=session,
)
assert get_mapped_ti_count(consume_task, dr.run_id, session=session) == 2
+
+
[email protected](
+ ("operation", "expected_length"),
+ [("plain", 2), ("map", 2), ("zip", 2), ("zip_longest", 5), ("concat", 7)],
+)
+def test_map_length_selects_retained_producer_in_callers_iteration(
+ dag_maker, session, operation, expected_length
+):
+ with dag_maker(session=session, serialized=True) as dag:
+
+ @dag.task
+ def source():
+ return [1, 2]
+
+ @dag.task
+ def outside():
+ return [1, 2, 3, 4, 5]
+
+ source()
+ outside()
+
+ dr = dag_maker.create_dagrun()
+ original = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id,
node_id="loop")
+ session.add(original)
+ session.flush()
+ replacement = DynamicRegion(
+ dag_id=dr.dag_id,
+ run_id=dr.run_id,
+ node_id="loop",
+ forked_from_region_id=original.id,
+ resumes_from_index=1,
+ )
+ session.add(replacement)
+ source_ti = next(ti for ti in dr.task_instances if ti.task_id == "source")
+ outside_ti = next(ti for ti in dr.task_instances if ti.task_id ==
"outside")
+ source_ti.region_id = original.id
+ source_ti.region_index = 1
+ source_ti.state = TaskInstanceState.SUCCESS
+ earlier_ti = TaskInstance(
+ task=dag_maker.serialized_dag.get_task("source"),
+ run_id=dr.run_id,
+ map_index=0,
+ dag_version_id=source_ti.dag_version_id,
+ )
+ earlier_ti.region_id = original.id
+ earlier_ti.state = TaskInstanceState.SUCCESS
+ session.add(earlier_ti)
+ session.flush()
+ for ti, length in ((earlier_ti, 99), (source_ti, 2), (outside_ti, 5)):
+ XComModel.set_for_attempt(
+ task_instance_id=ti.id,
+ key=XCOM_RETURN_KEY,
+ value=list(range(length)),
+ mapped_length=length,
+ session=session,
+ )
+ session.flush()
+ source_arg =
SchedulerPlainXComArg(dag_maker.serialized_dag.get_task("source"),
XCOM_RETURN_KEY)
+ outside_arg =
SchedulerPlainXComArg(dag_maker.serialized_dag.get_task("outside"),
XCOM_RETURN_KEY)
+ argument = {
+ "plain": source_arg,
+ "map": SchedulerMapXComArg(source_arg, ["str"]),
+ "zip": SchedulerZipXComArg([source_arg, outside_arg], NOTSET),
+ "zip_longest": SchedulerZipXComArg([source_arg, outside_arg], None),
+ "concat": SchedulerConcatXComArg([source_arg, outside_arg]),
+ }[operation]
+ contexts = {"source": ProducerContext(replacement.id, 1,
loop_node_id="loop")}
+
+ assert (
+ get_task_map_length(argument, dr.run_id, producer_contexts=contexts,
session=session)
+ == expected_length
+ )
+
+
[email protected](
+ ("count", "unfinished", "expected_length"), [(0, False, 0), (2, False, 2),
(2, True, None)]
+)
+def test_mapped_producer_length_ignores_other_iterations(
+ dag_maker, session, count, unfinished, expected_length
+):
+ with dag_maker(session=session, serialized=True) as dag:
+
+ @dag.task
+ def source(value):
+ return value
+
+ source.expand(value=list(range(count)))
+
+ dr = dag_maker.create_dagrun()
+ loop = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id, node_id="loop")
+ session.add(loop)
+ session.flush()
+ expansions = [
+ DynamicRegion(
+ dag_id=dr.dag_id,
+ run_id=dr.run_id,
+ node_id="source",
+ parent_region_id=loop.id,
+ parent_region_index=index,
+ )
+ for index in range(2)
+ ]
+ session.add_all(expansions)
+ session.flush()
+ for ti in dr.task_instances:
+ ti.region_id = expansions[1].id
+ ti.state = TaskInstanceState.SUCCESS if count else
TaskInstanceState.SKIPPED
+ session.flush()
+ for ti in dr.task_instances:
+ if ti.map_index >= 0:
+ XComModel.set_for_attempt(
+ task_instance_id=ti.id, key=XCOM_RETURN_KEY,
value=ti.map_index, session=session
+ )
+ if unfinished:
+ dr.task_instances[0].state = TaskInstanceState.RUNNING
+ earlier_ti = TaskInstance(
+ task=dag_maker.serialized_dag.get_task("source"),
+ run_id=dr.run_id,
+ map_index=0,
+ dag_version_id=dr.task_instances[0].dag_version_id,
+ )
+ earlier_ti.region_id = expansions[0].id
+ earlier_ti.state = TaskInstanceState.RUNNING
+ session.add(earlier_ti)
+ session.flush()
+ XComModel.set_for_attempt(
+ task_instance_id=earlier_ti.id, key=XCOM_RETURN_KEY, value="other
iteration", session=session
+ )
+ session.flush()
+ argument =
SchedulerPlainXComArg(dag_maker.serialized_dag.get_task("source"),
XCOM_RETURN_KEY)
+
+ assert (
+ get_task_map_length(
+ argument,
+ dr.run_id,
+ producer_contexts={"source": ProducerContext(loop.id, 1,
loop_node_id="loop")},
+ session=session,
+ )
+ == expected_length
+ )
+
+
+def test_member_of_mapped_task_group_has_no_map_length(dag_maker, session):
+ from airflow.sdk import task_group
Review Comment:
Can `from airflow.sdk import task_group` move to the top of the file?
There's no import cycle here.
--
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]