| 272 | ) |
| 273 | |
| 274 | def wrapped(*args, **kwargs): |
| 275 | from distributed.worker import get_worker |
| 276 | |
| 277 | worker = get_worker() |
| 278 | |
| 279 | env = get_worker().plugins.get("env", {}) |
| 280 | if isinstance(env, DictReturnWorkerPlugin): |
| 281 | env = env.setup_fn() |
| 282 | worker.plugins["env"] = env |
| 283 | |
| 284 | kw_keys = set(kwargs.keys()) |
| 285 | considered = kw_keys if accepts_var_kw else (param_names & kw_keys) |
| 286 | conflict = set(env.keys()) & considered |
| 287 | if conflict: |
| 288 | raise ValueError( |
| 289 | f"Argument conflict: {conflict} passed via both kwargs and env" |
| 290 | ) |
| 291 | |
| 292 | filtered_env = { |
| 293 | k: v for k, v in env.items() if (k in param_names) and (k not in kwargs) |
| 294 | } |
| 295 | return func(*args, **kwargs, **filtered_env) |
| 296 | |
| 297 | return wrapped |
| 298 | |