"""External payload for the loosely coupled execution target. This script is deliberately NOT a Dagster code location. It depends only on ``dagster-pipes``, never on ``dagster``, and it never opens a connection to the orchestration metadata database. Everything it reports reaches the control plane through the pipes message channel. Run by ``distributed_execution.loosely_coupled.jobs`` via either ``PipesSubprocessClient`` (local) or ``PipesK8sClient`` (cluster). """ import os import socket from dagster_pipes import open_dagster_pipes # Presence of any of these would mean the payload was granted orchestration # runtime credentials it has no business holding. ORCHESTRATION_ENV_VARS = ( "DAGSTER_POSTGRES_HOST", "DAGSTER_POSTGRES_USER", "DAGSTER_POSTGRES_DB", ) def main() -> None: with open_dagster_pipes() as pipes: units = pipes.get_extra("units") host = socket.gethostname() worker = f"{host}#{os.getpid()}" pipes.log.info(f"External payload started on {worker} with {len(units)} work units") results = [ {"unit": unit, "squared": unit * unit, "host": host, "worker": worker} for unit in units ] leaked = [name for name in ORCHESTRATION_ENV_VARS if os.environ.get(name)] if leaked: pipes.log.warning( "Orchestration runtime credentials are visible to this payload: " f"{', '.join(leaked)}. A loosely coupled target should not have them." ) # The only channel back to the control plane. pipes.report_custom_message( { "results": results, "host": host, "worker": worker, "orchestration_env_visible": leaked, } ) pipes.log.info(f"External payload finished on {worker}") if __name__ == "__main__": main()