From 51dabb63e9010513305233921f97aab44a0dd02b Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 11 Sep 2026 11:24:06 +0000 Subject: [PATCH 1/2] Initial plan From b72f23a3deb17434f105fef0d9057e66839324c5 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 11 Sep 2026 11:34:54 +0000 Subject: [PATCH 2/2] Add mutable runtime variables for recipe execution Co-authored-by: ebhills <53243273+ebhills@users.noreply.github.com> --- README.md | 41 +++ schema/generate_recipe_schema.py | 2 +- schema/recipe_base_schema.json | 4 + tests/recipes/test_run.py | 361 +++++++++++++++++++++++ wrangles/connectors/__init__.py | 1 + wrangles/connectors/recipe.py | 124 +++++++- wrangles/connectors/variables.py | 58 ++++ wrangles/recipe.py | 484 ++++++++++++++++++++++++++++--- wrangles/recipe_wrangles/main.py | 24 +- wrangles/utils.py | 10 +- 10 files changed, 1068 insertions(+), 41 deletions(-) create mode 100644 wrangles/connectors/variables.py diff --git a/README.md b/README.md index 22cc4f1e8..934117c3f 100644 --- a/README.md +++ b/README.md @@ -211,3 +211,44 @@ write: - file: name: file.xlsx ``` + +#### Mutable runtime variables + +Recipes can now create and mutate runtime variables during a single run. + +- Set/update values with the `variables` run action. +- Capture any run action return value with `result_variable`. +- Reference current runtime values at execution time with `${runtime.variable}` (or dotted paths like `${runtime.upload_ref.file_id}`). +- Escape a runtime template literal with `\${runtime.variable}`. + +```yaml +run: + on_start: + - custom.build_reference: + result_variable: upload_ref + - variables: + update: + upload_ref: + status: ready + inspect: + - upload_ref.file_id + - upload_ref.status + +read: + - test: + rows: 1 + values: + product: A100 + +wrangles: + - create.column: + output: reference_id + value: ${runtime.upload_ref.file_id} + - log: + runtime_variables: + - upload_ref.status + log_data: false +``` + +Runtime variables are isolated per `wrangles.recipe.run()` invocation. Nested recipes inherit a snapshot by default; to export selected runtime variables back to the parent recipe, set `export_runtime_variables` on the `recipe` connector. Parallel branches (`concurrent`, `matrix`) run with isolated runtime copies unless exported explicitly. +Static `${variable}` templates are still resolved when the recipe loads. Use `${runtime.variable}` for mutable values that must resolve immediately before each step. Protected names (`row_count`, `column_count`, `columns`, `df`, `recipe_variables`, `applied_permission_group`) cannot be reassigned. diff --git a/schema/generate_recipe_schema.py b/schema/generate_recipe_schema.py index ad929f09f..809e7d5ef 100644 --- a/schema/generate_recipe_schema.py +++ b/schema/generate_recipe_schema.py @@ -222,7 +222,7 @@ def getMethodDocs(schema_wrangles, obj, path): run_properties = schema['run'][run]['properties'] - for x in ["if"]: + for x in ["if", "result_variable"]: run_properties[x] = { "$ref": f"#/$defs/run/commonProperties/{x}" } diff --git a/schema/recipe_base_schema.json b/schema/recipe_base_schema.json index c3d196e53..c56f16c30 100644 --- a/schema/recipe_base_schema.json +++ b/schema/recipe_base_schema.json @@ -93,6 +93,10 @@ "if": { "type": "string", "description": "Specify a condition to determine if this will execute or not. e.g. ${variable} == 1. Recipe variables ${variable} are parameterized and may be used within the statement." + }, + "result_variable": { + "type": "string", + "description": "Save the action return value into a mutable runtime variable." } } }, diff --git a/tests/recipes/test_run.py b/tests/recipes/test_run.py index c50835246..28e922ab3 100644 --- a/tests/recipes/test_run.py +++ b/tests/recipes/test_run.py @@ -9,6 +9,7 @@ import pandas as pd import wrangles import pytest +from wrangles.connectors import memory def test_on_success(): @@ -183,3 +184,363 @@ def test_run_pure_posix_path(): """ df = wrangles.recipe.run(pathlib.PurePosixPath('tests/samples/recipe-basic.wrgl.yml')) assert list(df.columns) == ['header1', 'header2'] + + +def _make_reference(): + return { + "file_id": "file-123", + "meta": { + "bucket": "docs", + "tags": ["a", "b"] + }, + "nullable": None + } + + +def test_runtime_result_variable_and_late_resolution(): + """ + Action results can be captured and consumed via runtime references. + """ + df = wrangles.recipe.run( + """ + run: + on_start: + - custom._make_reference: + result_variable: upload_ref + wrangles: + - create.column: + output: file_id + value: ${runtime.upload_ref.file_id} + - create.column: + output: metadata + value: ${runtime.upload_ref.meta} + - create.column: + output: nullable + value: ${runtime.upload_ref.nullable} + """, + dataframe=pd.DataFrame({"row": [1]}), + functions=_make_reference + ) + assert df["file_id"][0] == "file-123" + assert df["metadata"][0] == {"bucket": "docs", "tags": ["a", "b"]} + assert df["nullable"][0] == "" + + +def test_runtime_set_update_and_condition(): + """ + Runtime variables can be set, updated and used in conditions. + """ + df = wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + runtime_state: + retry: 0 + run_mode: active + - variables: + update: + runtime_state: + retry: 1 + nested: + value: done + wrangles: + - create.column: + output: retry + value: ${runtime.runtime_state.retry} + if: ${runtime.run_mode} == "active" + - create.column: + output: nested + value: ${runtime.runtime_state.nested} + """, + dataframe=pd.DataFrame({"row": [1]}) + ) + assert df["retry"][0] == 1 + assert df["nested"][0] == {"value": "done"} + + +def test_runtime_reference_without_dataframe(): + """ + Runtime variables are usable even when recipe starts without a dataframe. + """ + captured = [] + + def capture_value(value): + captured.append(value) + + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + startup_value: ready + on_success: + - custom.capture_value: + value: ${runtime.startup_value} + """, + functions=capture_value + ) + assert captured == ["ready"] + + +def test_runtime_variable_assignment_protects_reserved_names(): + with pytest.raises(ValueError, match="protected"): + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + row_count: 10 + """ + ) + + +def test_runtime_missing_variable_and_escaped_literal(caplog): + with pytest.raises(ValueError, match="Runtime variable '\\$\\{runtime.missing\\}' was not found"): + wrangles.recipe.run( + """ + read: + - test: + rows: 1 + values: + header: value + wrangles: + - create.column: + output: out + value: ${runtime.missing} + """ + ) + + wrangles.recipe.run( + """ + read: + - test: + rows: 1 + values: + header: value + wrangles: + - log: + info: '\\${runtime.literal_value}' + log_data: false + """ + ) + assert "${runtime.literal_value}" in caplog.messages + + +def test_runtime_nested_export_and_conflict_policy(): + """ + Nested recipes receive a snapshot and only export back explicitly. + """ + df = wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + parent_flag: parent + - recipe: + run: + on_start: + - variables: + set: + parent_flag: child + child_export: exported + export_runtime_variables: + - child_export + wrangles: + - create.column: + output: parent_flag + value: ${runtime.parent_flag} + - create.column: + output: child_export + value: ${runtime.child_export} + """, + dataframe=pd.DataFrame({"row": [1]}) + ) + assert df["parent_flag"][0] == "parent" + assert df["child_export"][0] == "exported" + + with pytest.raises(ValueError, match="Runtime export collision"): + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + parent_flag: parent + - recipe: + run: + on_start: + - variables: + set: + parent_flag: child + export_runtime_variables: + - parent_flag + """ + ) + + df_overwrite = wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + parent_flag: parent + - recipe: + run: + on_start: + - variables: + set: + parent_flag: child + export_runtime_variables: + - parent_flag + export_conflict_policy: overwrite + wrangles: + - create.column: + output: parent_flag + value: ${runtime.parent_flag} + """, + dataframe=pd.DataFrame({"row": [1]}) + ) + assert df_overwrite["parent_flag"][0] == "child" + + +def test_runtime_inspection_redacts_sensitive_values_and_bounds_output(caplog): + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + access_token: abc123 + signed_url: https://example.com/file?X-Amz-Signature=secret&X-Amz-Credential=cred + safe_info: + rows: + - 1 + - 2 + - 3 + - 4 + mark_secret: + - signed_url + inspect: + - access_token + - signed_url + - safe_info + max_items: 2 + """ + ) + message = "\n".join(caplog.messages) + assert "[REDACTED]" in message + assert "2 more item(s)" in message + + +def test_log_wrangle_runtime_variables(caplog): + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + upload_ref: + file_id: file-123 + read: + - test: + rows: 1 + values: + col1: value + wrangles: + - log: + runtime_variables: + - upload_ref.file_id + log_data: false + """ + ) + + assert "upload_ref.file_id" in "\n".join(caplog.messages) + + +def test_runtime_variables_do_not_mutate_caller_input_dict(): + variables = { + "mutable_settings": { + "value": "original" + } + } + + wrangles.recipe.run( + """ + run: + on_start: + - variables: + update: + mutable_settings: + value: changed + """, + variables=variables + ) + + assert variables["mutable_settings"]["value"] == "original" + + +def test_runtime_parallel_branches_are_isolated(): + observed = [] + + def capture_branch(value): + observed.append(value) + + wrangles.recipe.run( + """ + run: + on_start: + - concurrent: + run: + - recipe: + run: + on_start: + - variables: + set: + branch_value: 1 + - custom.capture_branch: + value: ${runtime.branch_value} + - recipe: + run: + on_start: + - variables: + set: + branch_value: 2 + - custom.capture_branch: + value: ${runtime.branch_value} + """, + functions=capture_branch + ) + + assert sorted(observed) == [1, 2] + + +def test_runtime_references_in_read_write_and_if(): + memory.clear() + memory.dataframes["runtime_source"] = pd.DataFrame({"value": ["ok"]}) + + wrangles.recipe.run( + """ + run: + on_start: + - variables: + set: + source_id: runtime_source + target_id: runtime_target + should_write: true + read: + - memory: + id: ${runtime.source_id} + write: + - memory: + id: ${runtime.target_id} + if: ${runtime.should_write} == True + """ + ) + + assert "runtime_target" in memory.dataframes + assert memory.dataframes["runtime_target"]["data"][0][0] == "ok" diff --git a/wrangles/connectors/__init__.py b/wrangles/connectors/__init__.py index 1c50c4c0e..bd0d20f65 100644 --- a/wrangles/connectors/__init__.py +++ b/wrangles/connectors/__init__.py @@ -27,5 +27,6 @@ from . import s3 from . import train from . import jinja +from . import variables from . import _formatting from . import input diff --git a/wrangles/connectors/recipe.py b/wrangles/connectors/recipe.py index 711818482..4af8aad79 100644 --- a/wrangles/connectors/recipe.py +++ b/wrangles/connectors/recipe.py @@ -17,6 +17,8 @@ def run( name: str = None, variables: dict = None, functions: _Union[_types.FunctionType, list] = [], + export_runtime_variables: _Union[str, list] = None, + export_conflict_policy: str = "error", **kwargs ) -> None: """ @@ -28,11 +30,39 @@ def run( :param name: Name of the recipe to run :param variables: (Optional) A dictionary of custom variables to override placeholders in the recipe. Variables can be indicated as ${MY_VARIABLE}. Variables can also be overwritten by Environment Variables. :param functions: Pass in a custom function or list of custom functions that can be called in the recipe. + :param export_runtime_variables: Runtime variable name(s) to export back to the parent recipe. + :param export_conflict_policy: How to resolve parent collisions when exporting runtime variables. """ if variables is None: variables = {} if not name: name = kwargs - _recipe.run(name, variables=variables, functions=functions) + if export_runtime_variables is None: + _recipe.run(name, variables=variables, functions=functions) + return + + _, exported = _recipe.run( + name, + variables=variables, + functions=functions, + _return_runtime_variables=export_runtime_variables + ) + + if export_conflict_policy not in ["error", "overwrite"]: + raise ValueError("export_conflict_policy must be 'error' or 'overwrite'.") + + if export_conflict_policy == "error": + collisions = [ + key + for key, value in exported.items() + if key in variables and variables[key] != value + ] + if collisions: + raise ValueError( + f"Runtime export collision for variables {collisions}. " + "Use export_conflict_policy: overwrite to allow replacing parent values." + ) + + _recipe._apply_runtime_variable_updates(variables, set_values=exported) _schema['run'] = """ anyOf: @@ -47,6 +77,17 @@ def run( variables: type: object description: A dictionary of variables to pass to the recipe + export_runtime_variables: + type: + - string + - array + description: Runtime variable names to export from the nested recipe back to the parent runtime variables. + export_conflict_policy: + type: string + enum: + - error + - overwrite + description: Collision policy when exporting runtime variables back to the parent recipe. """ @@ -55,6 +96,8 @@ def read( variables: dict = None, columns: list = None, functions: _Union[_types.FunctionType, list] = [], + export_runtime_variables: _Union[str, list] = None, + export_conflict_policy: str = "error", **kwargs ) -> _pd.DataFrame: """ @@ -67,11 +110,35 @@ def read( :param variables: (Optional) A dictionary of custom variables to override placeholders in the recipe. Variables can be indicated as ${MY_VARIABLE}. Variables can also be overwritten by Environment Variables. :param columns: (Optional) Subset of the columns to include from the output of the recipe. If not provided, all columns will be included. :param functions: Pass in a custom function or list of custom functions that can be called in the recipe. + :param export_runtime_variables: Runtime variable name(s) to export back to the parent recipe. + :param export_conflict_policy: How to resolve parent collisions when exporting runtime variables. """ if variables is None: variables = {} if not name: name = kwargs - df = _recipe.run(name, variables=variables, functions=functions) + if export_runtime_variables is None: + df = _recipe.run(name, variables=variables, functions=functions) + else: + df, exported = _recipe.run( + name, + variables=variables, + functions=functions, + _return_runtime_variables=export_runtime_variables + ) + if export_conflict_policy not in ["error", "overwrite"]: + raise ValueError("export_conflict_policy must be 'error' or 'overwrite'.") + if export_conflict_policy == "error": + collisions = [ + key + for key, value in exported.items() + if key in variables and variables[key] != value + ] + if collisions: + raise ValueError( + f"Runtime export collision for variables {collisions}. " + "Use export_conflict_policy: overwrite to allow replacing parent values." + ) + _recipe._apply_runtime_variable_updates(variables, set_values=exported) # Select only specific columns if user requests them if columns is not None: @@ -99,6 +166,17 @@ def read( description: >- Subset of the columns to include from the output of the recipe. If not provided, all columns will be included. + export_runtime_variables: + type: + - string + - array + description: Runtime variable names to export from the nested recipe back to the parent runtime variables. + export_conflict_policy: + type: string + enum: + - error + - overwrite + description: Collision policy when exporting runtime variables back to the parent recipe. """ @@ -108,6 +186,8 @@ def write( variables: dict = None, columns: list = None, functions: _Union[_types.FunctionType, list] = [], + export_runtime_variables: _Union[str, list] = None, + export_conflict_policy: str = "error", **kwargs ) -> None: """ @@ -121,6 +201,8 @@ def write( :param variables: (Optional) A dictionary of custom variables to override placeholders in the recipe. Variables can be indicated as ${MY_VARIABLE}. Variables can also be overwritten by Environment Variables. :param columns: (Optional) A list of the columns to pass to the recipe. If omitted, all columns will be included. :param functions: Pass in a custom function or list of custom functions that can be called in the recipe. + :param export_runtime_variables: Runtime variable name(s) to export back to the parent recipe. + :param export_conflict_policy: How to resolve parent collisions when exporting runtime variables. """ if variables is None: variables = {} @@ -130,7 +212,32 @@ def write( columns = _wildcard_expansion(df.columns, columns) df = df[columns] - _recipe.run(name, dataframe=df, variables=variables, functions=functions) + if export_runtime_variables is None: + _recipe.run(name, dataframe=df, variables=variables, functions=functions) + return + + _, exported = _recipe.run( + name, + dataframe=df, + variables=variables, + functions=functions, + _return_runtime_variables=export_runtime_variables + ) + + if export_conflict_policy not in ["error", "overwrite"]: + raise ValueError("export_conflict_policy must be 'error' or 'overwrite'.") + if export_conflict_policy == "error": + collisions = [ + key + for key, value in exported.items() + if key in variables and variables[key] != value + ] + if collisions: + raise ValueError( + f"Runtime export collision for variables {collisions}. " + "Use export_conflict_policy: overwrite to allow replacing parent values." + ) + _recipe._apply_runtime_variable_updates(variables, set_values=exported) _schema['write'] = """ @@ -146,4 +253,15 @@ def write( variables: type: object description: A dictionary of variables to pass to the recipe + export_runtime_variables: + type: + - string + - array + description: Runtime variable names to export from the nested recipe back to the parent runtime variables. + export_conflict_policy: + type: string + enum: + - error + - overwrite + description: Collision policy when exporting runtime variables back to the parent recipe. """ diff --git a/wrangles/connectors/variables.py b/wrangles/connectors/variables.py new file mode 100644 index 000000000..ac464b7aa --- /dev/null +++ b/wrangles/connectors/variables.py @@ -0,0 +1,58 @@ +""" +Manage mutable runtime recipe variables. +""" +import logging as _logging +from typing import Union as _Union +from .. import recipe as _recipe + + +_schema = {} + + +def run( + variables: dict = None, + inspect: _Union[str, list] = None, + max_items: int = 20 +): + """ + type: object + description: Manage mutable runtime variables during recipe execution. + properties: + set: + type: object + description: Set one or more runtime variables atomically. + update: + type: object + description: Deep-merge dictionary patches into existing runtime variables. + mark_secret: + type: + - string + - array + description: Mark runtime variable names as sensitive for redaction. + inspect: + type: + - string + - array + description: Variable name(s) or dotted paths to log in redacted form. + max_items: + type: integer + minimum: 1 + description: Maximum list/dictionary entries to include during inspection. + """ + if variables is None: + variables = {} + + if inspect: + snapshot = _recipe.inspect_runtime_variables( + variables=variables, + include=inspect, + max_items=max_items + ) + if snapshot: + _logging.info(f": Runtime Variables :: {snapshot}") + return snapshot + + return None + + +_schema['run'] = run.__doc__ diff --git a/wrangles/recipe.py b/wrangles/recipe.py index 934c20a53..8687a4b7f 100644 --- a/wrangles/recipe.py +++ b/wrangles/recipe.py @@ -17,6 +17,7 @@ import time as _time import pandas as _pandas import requests as _requests +from urllib.parse import urlsplit as _urlsplit, parse_qsl as _parse_qsl, urlunsplit as _urlunsplit from . import recipe_wrangles as _recipe_wrangles from . import connectors as _connectors from . import data as _data @@ -52,6 +53,337 @@ default=None ) +_RUNTIME_SECRET_VARIABLES_KEY = "__runtime_secret_variables__" +_RUNTIME_INTERNAL_KEYS = {_RUNTIME_SECRET_VARIABLES_KEY} +_RUNTIME_PROTECTED_NAMES = { + "recipe_variables", + "row_count", + "column_count", + "columns", + "df", + "applied_permission_group", + *_RUNTIME_INTERNAL_KEYS +} +_RUNTIME_REFERENCE_PATTERN = _re.compile(r"(? dict: + return { + key: value + for key, value in variables.items() + if key not in _RUNTIME_INTERNAL_KEYS and key != "recipe_variables" + } + + +def _refresh_recipe_variables(variables: dict) -> None: + variables["recipe_variables"] = _clone_runtime_value(_runtime_assignable_variables(variables)) + + +def _build_runtime_variables_view( + variables: dict, + df: _pandas.DataFrame = None +) -> dict: + variable_view = _clone_runtime_value(_runtime_assignable_variables(variables)) + + if isinstance(df, _pandas.DataFrame): + variable_view.update({ + "row_count": len(df), + "column_count": len(df.columns), + "columns": df.columns.tolist(), + "df": df + }) + + variable_view["recipe_variables"] = { + key: value + for key, value in variable_view.items() + if key != "recipe_variables" + } + + return variable_view + + +def _resolve_runtime_path(path: str, variables: dict): + current = variables + for segment in path.split("."): + if not isinstance(current, dict) or segment not in current: + raise ValueError(f"Runtime variable '${{runtime.{path}}}' was not found.") + current = current[segment] + return current + + +def _resolve_runtime_references( + recipe_object: _typing.Any, + variables: dict, + defer_keys: set = None +) -> _typing.Any: + if defer_keys is None: + defer_keys = set() + + if isinstance(recipe_object, list): + return [ + _resolve_runtime_references(item, variables, defer_keys) + for item in recipe_object + ] + + if isinstance(recipe_object, dict): + resolved = {} + for key, value in recipe_object.items(): + resolved_key = _resolve_runtime_references(key, variables, defer_keys) + if isinstance(resolved_key, str) and resolved_key in defer_keys: + resolved[resolved_key] = value + else: + resolved[resolved_key] = _resolve_runtime_references(value, variables, defer_keys) + return resolved + + if isinstance(recipe_object, str): + escaped_marker = "__WRANGLES_ESCAPED_RUNTIME_TEMPLATE__" + templated = recipe_object.replace(r"\${runtime.", f"{escaped_marker}") + + if _RUNTIME_REFERENCE_FULL_PATTERN.fullmatch(templated): + runtime_key = _RUNTIME_REFERENCE_FULL_PATTERN.fullmatch(templated).group(1) + return _resolve_runtime_path(runtime_key, variables) + + def _replace_runtime_variable(match): + runtime_key = match.group(1) + return str(_resolve_runtime_path(runtime_key, variables)) + + templated = _RUNTIME_REFERENCE_PATTERN.sub(_replace_runtime_variable, templated) + return templated.replace(escaped_marker, "${runtime.") + + return recipe_object + + +def _resolve_runtime_condition(statement: str, variables: dict) -> str: + if not isinstance(statement, str): + return statement + + escaped_marker = "__WRANGLES_ESCAPED_RUNTIME_TEMPLATE__" + statement = statement.replace(r"\${runtime.", f"{escaped_marker}") + + def _replace_runtime_variable(match): + runtime_key = match.group(1) + return repr(_resolve_runtime_path(runtime_key, variables)) + + statement = _RUNTIME_REFERENCE_PATTERN.sub(_replace_runtime_variable, statement) + return statement.replace(escaped_marker, "${runtime.") + + +def _validate_runtime_variable_name(name: str) -> None: + if not isinstance(name, str) or not _re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", name): + raise ValueError( + f"Runtime variable name '{name}' is invalid. " + "Use letters, numbers and underscores, starting with a letter or underscore." + ) + if name in _RUNTIME_PROTECTED_NAMES: + raise ValueError(f"Runtime variable '{name}' is protected and cannot be assigned.") + + +def _deep_merge_dicts(original: dict, updates: dict) -> dict: + merged = _clone_runtime_value(original) + for key, value in updates.items(): + if isinstance(merged.get(key), dict) and isinstance(value, dict): + merged[key] = _deep_merge_dicts(merged[key], value) + else: + merged[key] = _clone_runtime_value(value) + return merged + + +def _apply_runtime_variable_updates( + variables: dict, + set_values: dict = None, + update_values: dict = None, + mark_secret: list = None +) -> None: + prepared_assignments = {} + + if set_values is not None: + if not isinstance(set_values, dict): + raise ValueError("set must be a dictionary of runtime variable assignments.") + for key, value in set_values.items(): + _validate_runtime_variable_name(key) + prepared_assignments[key] = _clone_runtime_value(value) + + if update_values is not None: + if not isinstance(update_values, dict): + raise ValueError("update must be a dictionary of runtime variable updates.") + for key, patch in update_values.items(): + _validate_runtime_variable_name(key) + if key not in variables: + raise ValueError(f"Runtime variable '{key}' was not found for update.") + if not isinstance(variables[key], dict) or not isinstance(patch, dict): + raise ValueError( + f"Runtime variable '{key}' can only be updated with a dictionary patch." + ) + prepared_assignments[key] = _deep_merge_dicts(variables[key], patch) + + prepared_mark_secret = set() + if mark_secret is not None: + if isinstance(mark_secret, str): + mark_secret = [mark_secret] + if not isinstance(mark_secret, list) or not all(isinstance(v, str) for v in mark_secret): + raise ValueError("mark_secret must be a string or list of strings.") + for variable_name in mark_secret: + _validate_runtime_variable_name(variable_name) + prepared_mark_secret.add(variable_name) + + if prepared_assignments: + variables.update(prepared_assignments) + _refresh_recipe_variables(variables) + + if prepared_mark_secret: + existing = variables.get(_RUNTIME_SECRET_VARIABLES_KEY, set()) + if not isinstance(existing, set): + existing = set(existing) if isinstance(existing, list) else set() + variables[_RUNTIME_SECRET_VARIABLES_KEY] = existing.union(prepared_mark_secret) + + +def _looks_like_sensitive_name(name: str) -> bool: + if not isinstance(name, str): + return False + lowered = name.lower() + return any( + key in lowered + for key in [ + "password", "secret", "token", "api_key", "apikey", "access_key", + "authorization", "credential", "signature", "signed", "private_key" + ] + ) + + +def _redact_string(value: str) -> str: + if not isinstance(value, str): + return value + try: + parsed = _urlsplit(value) + if parsed.scheme and parsed.netloc and parsed.query: + sensitive_query_keys = { + "x-amz-signature", "x-amz-credential", "x-amz-security-token", + "x-amz-algorithm", "x-amz-date", "signature", "sig", "token", + "access_key", "awsaccesskeyid" + } + query_keys = {key.lower() for key, _ in _parse_qsl(parsed.query, keep_blank_values=True)} + if query_keys.intersection(sensitive_query_keys): + return _urlunsplit((parsed.scheme, parsed.netloc, parsed.path, "[REDACTED]", parsed.fragment)) + except Exception: + pass + return value + + +def _sanitize_runtime_value( + value, + *, + key_hint: str = None, + max_depth: int = 4, + max_items: int = 20, + max_string_length: int = 400, + secret_variables: set = None, + depth: int = 0 +): + if secret_variables is None: + secret_variables = set() + + if key_hint in secret_variables or _looks_like_sensitive_name(key_hint): + return "[REDACTED]" + + if depth >= max_depth: + return "" + + if isinstance(value, _pandas.DataFrame): + return f"" + if isinstance(value, (bytes, bytearray)): + return f"" + if isinstance(value, str): + value = _redact_string(value) + if len(value) > max_string_length: + return value[:max_string_length] + "... " + return value + if isinstance(value, dict): + keys = list(value.keys()) + sanitized = {} + for key in keys[:max_items]: + sanitized[key] = _sanitize_runtime_value( + value[key], + key_hint=str(key), + max_depth=max_depth, + max_items=max_items, + max_string_length=max_string_length, + secret_variables=secret_variables, + depth=depth + 1 + ) + if len(keys) > max_items: + sanitized["..."] = f"{len(keys) - max_items} more item(s)" + return sanitized + if isinstance(value, list): + sanitized = [ + _sanitize_runtime_value( + item, + key_hint=key_hint, + max_depth=max_depth, + max_items=max_items, + max_string_length=max_string_length, + secret_variables=secret_variables, + depth=depth + 1 + ) + for item in value[:max_items] + ] + if len(value) > max_items: + sanitized.append(f"... {len(value) - max_items} more item(s)") + return sanitized + return value + + +def inspect_runtime_variables( + variables: dict, + include: _typing.Union[list, str, None] = None, + df: _pandas.DataFrame = None, + max_items: int = 20 +) -> dict: + variable_view = _build_runtime_variables_view(variables, df) + secret_variables = variables.get(_RUNTIME_SECRET_VARIABLES_KEY, set()) + if not isinstance(secret_variables, set): + secret_variables = set(secret_variables) if isinstance(secret_variables, list) else set() + + if include is None: + return {} + if isinstance(include, str): + include = [include] + if not isinstance(include, list): + raise ValueError("runtime_variables must be a string or list of strings.") + + selected = {} + for selector in include: + if not isinstance(selector, str): + raise ValueError("runtime_variables entries must be strings.") + if selector == "recipe_variables": + value = variable_view.get("recipe_variables") + else: + value = _resolve_runtime_path(selector, variable_view) + selected[selector] = _sanitize_runtime_value( + value, + key_hint=selector.split(".")[0], + max_items=max_items, + secret_variables=secret_variables + ) + return selected + # Suppress pandas performance warnings # this appears in some instances during the recipe execution when generating new columns. @@ -258,12 +590,9 @@ def _load_recipe( recipe_object = _yaml.safe_load(recipe_string) - # Add variables to variables - variables['recipe_variables'] = { - key: value - for key, value in variables.items() - if key != 'recipe_variables' - } + if _RUNTIME_SECRET_VARIABLES_KEY not in variables: + variables[_RUNTIME_SECRET_VARIABLES_KEY] = set() + _refresh_recipe_variables(variables) # Keep a copy of the raw recipe string for line lookups. # @@ -307,7 +636,10 @@ def _load_recipe( break # Check if there are any templated valued to update - recipe_object = _replace_templated_values(recipe_object, variables) + recipe_object = _replace_templated_values( + recipe_object, + _build_runtime_variables_view(variables) + ) return recipe_object, functions @@ -417,12 +749,33 @@ def _run_actions( for action_type, params in action.items(): try: + if params is None: + params = {} + params = _resolve_runtime_references( + _clone_runtime_value(params), + _build_runtime_variables_view(variables), + _RUNTIME_DEFER_KEYS + ) + # If the action is conditional, check if it should be run if ( "if" in params and - not _evaluate_conditional(params["if"], variables) + not _evaluate_conditional( + _resolve_runtime_condition( + params["if"], + _build_runtime_variables_view(variables) + ), + _build_runtime_variables_view(variables) + ) ): continue + + result_variable = params.pop("result_variable", None) + set_values = params.pop("set", None) if action_type == "variables" else None + update_values = params.pop("update", None) if action_type == "variables" else None + mark_secret = params.pop("mark_secret", None) if action_type == "variables" else None + inspect_values = params.pop("inspect", None) if action_type == "variables" else None + inspect_max_items = params.pop("max_items", 20) if action_type == "variables" else 20 common_params = {} # Add to common_params dict and remove from params @@ -430,6 +783,28 @@ def _run_actions( if key in params.keys(): common_params[key] = params.pop(key) + if action_type == "variables": + _apply_runtime_variable_updates( + variables, + set_values=set_values, + update_values=update_values, + mark_secret=mark_secret + ) + result = None + if inspect_values is not None: + result = inspect_runtime_variables( + variables=variables, + include=inspect_values, + max_items=inspect_max_items + ) + _logging.info(f": Runtime Variables :: {result}") + if result_variable is not None: + _apply_runtime_variable_updates( + variables, + set_values={result_variable: result} + ) + continue + func = _get_nested_function(action_type, _connectors, functions, 'run') if action_type == "matrix": params['variables'] = {**variables, **params['variables']} if 'variables' in params else variables @@ -439,7 +814,13 @@ def _run_actions( _validate_function_args(func, args, action_type) # Execute the function - func(**args) + result = func(**args) + + if result_variable is not None: + _apply_runtime_variable_updates( + variables, + set_values={result_variable: result} + ) except Exception as e: # Wrap with enhanced error information _wrap_and_raise('ACTION', action_type, None, e) @@ -476,10 +857,24 @@ def _read_data( for read_type, read_params in read.items(): try: + if read_params is None: + read_params = {} + read_params = _resolve_runtime_references( + _clone_runtime_value(read_params), + _build_runtime_variables_view(variables, input_dataframe), + _RUNTIME_DEFER_KEYS + ) + # If the action is conditional, check if it should be run if ( "if" in read_params and - not _evaluate_conditional(read_params["if"], variables) + not _evaluate_conditional( + _resolve_runtime_condition( + read_params["if"], + _build_runtime_variables_view(variables, input_dataframe) + ), + _build_runtime_variables_view(variables, input_dataframe) + ) ): return None @@ -598,7 +993,13 @@ def _execute_wrangles( for wrangle, params in step.items(): try: - if params is None: params = {} + if params is None: + params = {} + params = _resolve_runtime_references( + _clone_runtime_value(params), + _build_runtime_variables_view(variables, df), + _RUNTIME_DEFER_KEYS + ) # Replace any conflicting reserved words with a safe alternative wrangle = _reserved_word_replacements.get(wrangle, wrangle) @@ -614,16 +1015,11 @@ def _execute_wrangles( if ( "if" in params and not _evaluate_conditional( - params["if"], - { - **variables, - **{ - "row_count": len(df), - "column_count": len(df.columns), - "columns": df.columns.tolist(), - "df": df - } - } + _resolve_runtime_condition( + params["if"], + _build_runtime_variables_view(variables, df) + ), + _build_runtime_variables_view(variables, df) ) ): _logging.info(f": Wrangling :: {wrangle} skipped due to not passing the if statement.") @@ -1101,6 +1497,14 @@ def _write_data( for export_type, params in export.items(): try: + if params is None: + params = {} + params = _resolve_runtime_references( + _clone_runtime_value(params), + _build_runtime_variables_view(variables, df), + _RUNTIME_DEFER_KEYS + ) + # Filter the dataframe as requested before passing # to the desired write function df_temp = _filter_dataframe(df, **params) @@ -1109,16 +1513,11 @@ def _write_data( if ( "if" in params and not _evaluate_conditional( - params["if"], - { - **variables, - **{ - "row_count": len(df_temp), - "column_count": len(df_temp.columns), - "columns": df_temp.columns.tolist(), - "df": df_temp - } - } + _resolve_runtime_condition( + params["if"], + _build_runtime_variables_view(variables, df) + ), + _build_runtime_variables_view(variables, df_temp) ) ): continue @@ -1237,7 +1636,8 @@ def run( variables: dict = None, dataframe: _pandas.DataFrame = None, functions: _Union[_types.FunctionType, list, dict] = [], - timeout: float = None + timeout: float = None, + _return_runtime_variables: list = None ) -> _pandas.DataFrame: """ Execute a Wrangles Recipe. Recipes are written in YAML and allow @@ -1257,7 +1657,7 @@ def run( """ if variables is None: variables = {} - variables = variables.copy() + variables = _clone_runtime_value(variables) parent_context = _RECIPE_RUN_CONTEXT.get() run_context = { @@ -1287,7 +1687,23 @@ def run( dataframe, functions ) - return future.result(timeout) + result_df = future.result(timeout) + if _return_runtime_variables is not None: + if isinstance(_return_runtime_variables, str): + _return_runtime_variables = [_return_runtime_variables] + if not isinstance(_return_runtime_variables, list): + raise ValueError("_return_runtime_variables must be a string or list of strings.") + exported = {} + runtime_values = _runtime_assignable_variables(variables) + for variable_name in _return_runtime_variables: + _validate_runtime_variable_name(variable_name) + if variable_name not in runtime_values: + raise ValueError( + f"Runtime variable '{variable_name}' was requested for export but was not found." + ) + exported[variable_name] = _clone_runtime_value(runtime_values[variable_name]) + return result_df, exported + return result_df except _futures.TimeoutError as e: try: diff --git a/wrangles/recipe_wrangles/main.py b/wrangles/recipe_wrangles/main.py index 28cfbc9a3..cd40dbd45 100644 --- a/wrangles/recipe_wrangles/main.py +++ b/wrangles/recipe_wrangles/main.py @@ -844,6 +844,8 @@ def log( error: str = None, warning: str = None, info: str = None, + runtime_variables: _Union[str, list] = None, + runtime_max_items: int = 20, log_data: bool = None, **kwargs ): @@ -870,12 +872,21 @@ def log( info: type: string description: Log info to the console + runtime_variables: + type: + - string + - array + description: Runtime variable names or dotted paths to inspect and log in redacted form. + runtime_max_items: + type: integer + minimum: 1 + description: Maximum list/dictionary entries to include per inspected runtime variable. log_data: type: boolean description: Whether to log a sample of the contents of the dataframe. Default True if not logging to a write, error, warning or info. Default False otherwise. """ - variables = kwargs.pop('variables', {}) - variables = _delayed_variable_interpretation(df, variables) + runtime_values = kwargs.pop('variables', {}) + variables = _delayed_variable_interpretation(df, runtime_values) # Handle variable interpretation for all parameters columns, write, error, warning, info = [ @@ -928,6 +939,15 @@ def log( dataframe=df ) + if runtime_variables: + runtime_snapshot = _recipe.inspect_runtime_variables( + variables=runtime_values, + include=runtime_variables, + df=df, + max_items=runtime_max_items + ) + _logging.info(f": Runtime Variables :: {runtime_snapshot}") + if log_data == None and not any([error, warning, info, write]): log_data = True if log_data: diff --git a/wrangles/utils.py b/wrangles/utils.py index 8a9f3098d..ec2c3286a 100644 --- a/wrangles/utils.py +++ b/wrangles/utils.py @@ -525,7 +525,11 @@ def delayed_variable_interpretation(df, variables=None): variables = {} variables = { - **variables, + **{ + k: v + for k, v in variables.items() + if not str(k).startswith("__runtime_") + }, **{ "row_count": len(df), "column_count": len(df.columns), @@ -600,6 +604,8 @@ def replace_templated_values( # Whole string is a variable if variable_pattern.fullmatch(new_recipe_object): + if new_recipe_object.startswith("${runtime."): + return new_recipe_object try: replacement_value = variables[new_recipe_object[2:-1]] except: @@ -647,6 +653,8 @@ def replace_templated_values( # Since this is within a string, the type is forced to also be a string elif variable_pattern.search(new_recipe_object): for var in variable_pattern.findall(new_recipe_object): + if var.startswith("${runtime."): + continue try: replacement_value = variables[var[2:-1]] except: