Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
78 changes: 78 additions & 0 deletions src/flowx/bundler/dab_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -1543,6 +1543,80 @@ def visit(task: dict[str, Any]) -> None:
return neutralized


# NOTE: This is a package-time *safety net*, not the ideal fix. The real gap is upstream in the convert
# phase: when the translator resolves ``@activity('P').output`` / ``@variables('X')`` into a
# ``{{tasks.P.values.*}}`` reference, it stores the value in the IR but does not promote that data-flow
# reference into a control-flow ``Activity.depends_on`` edge (unlike the ForEach/Switch/IfCondition
# bridges, which self-wire their edge at synthesis time). So the edge is often missing by the time we
# emit YAML. Backfilling here keeps ``bundle deploy`` from being rejected; wiring the edge in the
# translator would make this a redundant belt-and-suspenders check. See discussion on issue #27.
Comment on lines +1546 to +1552

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we move this to the convert phase? It seems like that's the right place to implement this, the references are already in the IR with missing edges.

I think we can define something like this, then call it at (engine.py:196):

def backfill_task_value_dependencies(pipeline: Pipeline, warnings: list[str]) -> Pipeline:
    activity_indices = {activity.task_key: i for i, activity in enumerate(pipeline.tasks)}
    for activity in pipeline.tasks:
        for producer in task_value_refs(activity): 
            if producer in upstream_keys(activity):
                continue
            if activity_indices[producer] > activity_indices[activity.task_key]:
                warnings.append(
                    f"{activity.task_key} reads a value from '{producer}', which runs later"
                )
                continue
            activity.depends_on = [*(activity.depends_on or []), Dependency(task_key=producer)]
    return pipeline

def _backfill_task_value_dependencies(tasks: list[dict[str, Any]]) -> int:
"""Adds a ``depends_on`` edge for every ``{{tasks.P.values.*}}`` reference.

Databricks rejects a job whose task references ``{{tasks.P.values.Y}}`` when
``P`` is not a (transitive) dependency of the referencing task
(``INVALID_PARAMETER_VALUE: Task 'P' must be a dependency of task 'T'``).

flowx rewrites ADF ``@variables('X')`` / ``@activity('Y').output`` into such
task-value references (e.g. synthesised ``_init_<var>`` tasks and the
IfCondition/Switch/ForEach bridge tasks) but does not always add the matching
dependency edge. This pass walks every task's ``condition_task.{left,right}``
and ``notebook_task.base_parameters`` and, for each referenced producer that
exists in this job's task list but is not already a *direct* dependency of the
consumer, appends ``{"task_key": P}`` to the consumer's ``depends_on``.

Runs *after* :func:`_strip_dangling_task_value_refs`, so refs to producers
absent from the bundle have already been blanked and are not seen here.
Operates within a single job scope; ``for_each_task.task`` bodies are handled
against their own sibling set because inner-body task keys are local to the
parent task list. Returns the number of edges added.
"""
added = 0

def _referenced_producers(task: dict[str, Any]) -> set[str]:
producers: set[str] = set()

def scan(value: Any) -> None:
if isinstance(value, str):
producers.update(_TASK_VALUE_REF.findall(value))
elif isinstance(value, dict):
for item in value.values():
scan(item)
elif isinstance(value, list):
for item in value:
scan(item)

notebook_task = task.get("notebook_task") or {}
scan(notebook_task.get("base_parameters"))
condition_task = task.get("condition_task") or {}
scan(condition_task.get("left"))
scan(condition_task.get("right"))
return producers

def process(scope_tasks: list[dict[str, Any]]) -> None:
nonlocal added
sibling_keys = {t.get("task_key", "") for t in scope_tasks}
sibling_keys.discard("")
for task in scope_tasks:
self_key = task.get("task_key", "")
existing = {dep.get("task_key") for dep in (task.get("depends_on") or []) if isinstance(dep, dict)}
for producer in sorted(_referenced_producers(task)):
if producer == self_key or producer in existing or producer not in sibling_keys:
continue
deps = list(task.get("depends_on") or [])
deps.append({"task_key": producer})
task["depends_on"] = deps
existing.add(producer)
added += 1
# Recurse into for_each bodies against their own sibling scope.
for_each = task.get("for_each_task")
if for_each and isinstance(for_each.get("task"), dict):
process([for_each["task"]])

process(tasks)
return added


def _collect_all_task_keys(tasks: list[dict[str, Any]]) -> set[str]:
"""Collects every task_key reachable from the job's top-level task list."""
keys: set[str] = set()
Expand Down Expand Up @@ -1623,6 +1697,10 @@ def _build_job_resource(
_neutralized_conditions.extend(
_strip_dangling_task_value_refs(workflow.tasks, _collect_all_task_keys(workflow.tasks))
)
# A task-value producer must be a dependency of any task that references it, or
# Databricks rejects the job at create time. Backfill the missing edges for refs
# to producers that survived the dangling-ref strip above.
_backfill_task_value_dependencies(workflow.tasks)

job_def: dict[str, Any] = {
"name": workflow.name,
Expand Down
109 changes: 109 additions & 0 deletions tests/unit/test_bundler.py
Original file line number Diff line number Diff line change
Expand Up @@ -962,6 +962,115 @@ def test_recurses_into_for_each_task_body(self):
assert tasks[0]["for_each_task"]["task"]["run_job_task"]["job_parameters"]["x"] == ""


class TestBackfillTaskValueDependencies:
"""PR #35 (#30): a task that references ``{{tasks.P.values.*}}`` must have ``P`` as a
direct dependency, or Databricks rejects the job at create time
(``INVALID_PARAMETER_VALUE: Task 'P' must be a dependency of task 'T'``).
``_backfill_task_value_dependencies`` adds the missing ``depends_on`` edge."""

def test_condition_operand_ref_backfills_dependency(self):
# The core case: a condition operand reads a producer's task value but the edge
# was never wired, so deploy would be rejected. The backfill adds it.
from flowx.bundler.dab_writer import _backfill_task_value_dependencies

tasks = [
{"task_key": "_init_flag", "notebook_task": {"notebook_path": "/n"}},
{
"task_key": "gate",
"condition_task": {
"op": "EQUAL_TO",
"left": "{{tasks._init_flag.values.flag}}",
"right": "true",
},
},
]
added = _backfill_task_value_dependencies(tasks)
assert added == 1
assert tasks[1]["depends_on"] == [{"task_key": "_init_flag"}]

def test_notebook_base_parameter_ref_backfills_dependency(self):
from flowx.bundler.dab_writer import _backfill_task_value_dependencies

tasks = [
{"task_key": "producer", "notebook_task": {"notebook_path": "/p"}},
{
"task_key": "consumer",
"notebook_task": {
"notebook_path": "/c",
"base_parameters": {"run_id": "{{tasks.producer.values.run_id}}"},
},
},
]
added = _backfill_task_value_dependencies(tasks)
assert added == 1
assert tasks[1]["depends_on"] == [{"task_key": "producer"}]

def test_does_not_duplicate_an_existing_dependency(self):
from flowx.bundler.dab_writer import _backfill_task_value_dependencies

tasks = [
{"task_key": "producer", "notebook_task": {"notebook_path": "/p"}},
{
"task_key": "consumer",
"depends_on": [{"task_key": "producer"}],
"condition_task": {
"op": "EQUAL_TO",
"left": "{{tasks.producer.values.x}}",
"right": "1",
},
},
]
added = _backfill_task_value_dependencies(tasks)
assert added == 0
assert tasks[1]["depends_on"] == [{"task_key": "producer"}]

def test_skips_producer_absent_from_scope(self):
# A ref whose producer is not a sibling in this scope (e.g. blanked earlier by
# _strip_dangling_task_value_refs) must not fabricate an edge to a missing task.
from flowx.bundler.dab_writer import _backfill_task_value_dependencies

tasks = [
{
"task_key": "consumer",
"condition_task": {
"op": "EQUAL_TO",
"left": "{{tasks.ghost.values.x}}",
"right": "1",
},
},
]
added = _backfill_task_value_dependencies(tasks)
assert added == 0
assert "depends_on" not in tasks[0]

def test_for_each_body_is_scoped_locally(self):
# Inner-body task keys are local to the parent list, so a body ref to an outer
# producer is not backfilled at body level (it recurses without erroring).
from flowx.bundler.dab_writer import _backfill_task_value_dependencies

tasks = [
{"task_key": "outer", "notebook_task": {"notebook_path": "/o"}},
{
"task_key": "loop",
"depends_on": [{"task_key": "outer"}],
"for_each_task": {
"inputs": "[1, 2, 3]",
"task": {
"task_key": "loop_body",
"condition_task": {
"op": "EQUAL_TO",
"left": "{{tasks.outer.values.x}}",
"right": "1",
},
},
},
},
]
added = _backfill_task_value_dependencies(tasks)
assert added == 0
assert "depends_on" not in tasks[1]["for_each_task"]["task"]


class TestAggregatedReportPipelineParameters:
"""Change pipeline-parameters-and-variables-round-trip (P0): VAR-001."""

Expand Down
Loading