The reference graph processed every work unit inside a single op, so exactly one worker was ever reported regardless of executor. That made the guide's claim that k8s_job_executor yields several distinct contributing_hosts false, and left the reference implementations unable to demonstrate the executor choice at all. generate_work_units is now a DynamicOut and both graphs map over it, so one step is created per unit and the loosely coupled pattern dispatches one external workload per unit. Evidence is split into contributing_workers (host and pid, differs per process) and contributing_hosts (differs only across machines), because the previous single field could not distinguish multiprocess fan-out from no fan-out. Tests now assert the mapped step keys rather than a host count, since execute_in_process ignores executor_def and cannot prove executor behaviour on its own. Adds a Windows note: multiprocess_executor did not complete during authoring and left orphaned processes. Changelog: fixed
96 lines
3.0 KiB
Python
96 lines
3.0 KiB
Python
"""Verifies the loosely coupled pipes reference end to end via the subprocess transport."""
|
|
|
|
import pytest
|
|
from dagster import DagsterInvariantViolationError, Failure
|
|
|
|
from distributed_execution.loosely_coupled.jobs import (
|
|
LOOSELY_COUPLED_JOBS,
|
|
PAYLOAD_SCRIPT,
|
|
loosely_coupled_k8s_job,
|
|
loosely_coupled_subprocess_job,
|
|
)
|
|
from distributed_execution.repository import defs
|
|
|
|
|
|
def test_payload_script_exists():
|
|
from pathlib import Path
|
|
|
|
assert Path(PAYLOAD_SCRIPT).is_file()
|
|
|
|
|
|
def test_payload_does_not_import_dagster():
|
|
from pathlib import Path
|
|
|
|
source = Path(PAYLOAD_SCRIPT).read_text(encoding="utf-8")
|
|
assert "from dagster_pipes import" in source
|
|
assert "import dagster\n" not in source
|
|
assert "from dagster import" not in source
|
|
|
|
|
|
def test_both_transports_are_registered():
|
|
names = {job.name for job in defs.jobs}
|
|
assert {"loosely_coupled_subprocess_job", "loosely_coupled_k8s_job"} <= names
|
|
assert len(LOOSELY_COUPLED_JOBS) == 2
|
|
|
|
|
|
@pytest.mark.parametrize(
|
|
("job", "transport"),
|
|
[
|
|
(loosely_coupled_subprocess_job, "subprocess"),
|
|
(loosely_coupled_k8s_job, "k8s_pod_logs"),
|
|
],
|
|
)
|
|
def test_jobs_declare_their_execution_target(job, transport):
|
|
assert job.tags["execution_target"] == "loosely_coupled"
|
|
assert job.tags["transport"] == transport
|
|
|
|
|
|
def test_subprocess_job_runs_end_to_end():
|
|
result = loosely_coupled_subprocess_job.execute_in_process()
|
|
|
|
assert result.success
|
|
summary = result.output_for_node("summarise_results")
|
|
# Same shape as the tightly coupled reference: only the middle node changed.
|
|
assert summary["units"] == 4
|
|
assert summary["total"] == 14
|
|
|
|
|
|
def test_each_unit_is_dispatched_to_its_own_external_worker():
|
|
result = loosely_coupled_subprocess_job.execute_in_process()
|
|
|
|
mapped = result.output_for_node("dispatch_external_work_subprocess")
|
|
assert set(mapped) == {"unit_0", "unit_1", "unit_2", "unit_3"}
|
|
# Each dispatch is its own OS process, so workers differ even on a single host.
|
|
assert len({row["worker"] for row in mapped.values()}) == 4
|
|
|
|
|
|
def test_payload_reports_no_orchestration_credentials(monkeypatch):
|
|
monkeypatch.delenv("DAGSTER_POSTGRES_HOST", raising=False)
|
|
monkeypatch.delenv("DAGSTER_POSTGRES_USER", raising=False)
|
|
monkeypatch.delenv("DAGSTER_POSTGRES_DB", raising=False)
|
|
|
|
result = loosely_coupled_subprocess_job.execute_in_process()
|
|
mapped = result.output_for_node("dispatch_external_work_subprocess")
|
|
|
|
assert len(mapped) == 4
|
|
assert all(row["worker"] for row in mapped.values())
|
|
|
|
|
|
def test_silent_message_path_is_treated_as_failure():
|
|
from distributed_execution.loosely_coupled.jobs import _result_from_pipes
|
|
|
|
class _Silent:
|
|
def get_custom_messages(self):
|
|
return []
|
|
|
|
class _Ctx:
|
|
class log: # noqa: N801
|
|
@staticmethod
|
|
def warning(*_args, **_kwargs): ...
|
|
|
|
@staticmethod
|
|
def info(*_args, **_kwargs): ...
|
|
|
|
with pytest.raises(Failure, match="No pipes messages received"):
|
|
_result_from_pipes(_Ctx(), _Silent())
|