Skip to content
Closed
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
12 changes: 10 additions & 2 deletions src/deadline_worker_agent/api_models.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,13 @@ class ChunkIntParameter(TypedDict):


class BoolParameter(TypedDict):
bool: bool
# Transitional union: the service now sends string-typed booleans drawn from
# the Open Job Description specification's case-insensitive boolean
# vocabulary (e.g. "true"/"false", "yes"/"no", "on"/"off", "1"/"0"), but
# jobs created before the flip still carry native booleans in persisted
# parameters that are returned verbatim, so either form may arrive during
# and after rollout.
bool: str | bool


class RangeExprParameter(TypedDict):
Expand All @@ -109,7 +115,9 @@ class FloatListParameter(TypedDict):


class BoolListParameter(TypedDict):
boolList: list[bool]
# Transitional union: see BoolParameter. Elements may be string-typed
# booleans or native booleans during and after the rollout.
boolList: list[str | bool]


class IntListListParameter(TypedDict):
Expand Down
16 changes: 15 additions & 1 deletion src/deadline_worker_agent/scheduler/session_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -532,7 +532,21 @@ def dequeue(self) -> SessionActionDefinition | None:
task_id=task_id,
) from e
task_parameters_data: dict = action_definition.get("parameters", {})
task_parameters = parameters_from_api_response(task_parameters_data)
# Task parameters are decoded here, outside the step_details
# fetch above. A decode failure (e.g. a malformed boolean) must
# fail this one action gracefully, matching how step_details
# errors are handled, rather than escaping as a bare ValueError
# that would tear down the whole session.
try:
task_parameters = parameters_from_api_response(task_parameters_data)
except (ValueError, RuntimeError) as e:

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

The new try closes the ValueError hole but leaves the adjacent TypeError hole open, so the session-teardown failure mode this block exists to prevent is still reachable.

parameters_from_api_response dispatches on "string" in value, "bool" in value, etc. (job_details.py:111-171) without first checking that value is a dict. Task parameters are never run through _validate_job_parameters — task_parameters_data comes straight from the unvalidated AssignedSession payload at line 534 — so if the service ever sends a parameter whose value is not a container (e.g. {"p": null} or {"p": 5}), the membership test raises TypeError: argument of type NoneType is not iterable, not ValueError. That escapes this except (ValueError, RuntimeError), propagates out of dequeue(), and since Session._start_action (sessions/session.py:653-658) only catches SessionActionError, it tears down the whole session and cancels every remaining queued action.

That is precisely the concern the PR itself raises one file over: the boolList shape-check comment at job_details.py:155-157 says it exists so callers "catching only ValueError" do not miss a TypeError. But only that one boolList instance was patched, while the dispatch chain above it has the same exposure for every parameter type.

Two cheap ways to close the class rather than the single instance:

  • add TypeError here: except (ValueError, RuntimeError, TypeError) as e:, or
  • shape-check at the top of the parameters_from_api_response loop, e.g. if not isinstance(value, dict): raise ValueError(f"Parameter {name} -- expected a dict but got {value!r}"), which also yields a message naming the offending parameter.

The second seems preferable since it fixes both call sites and keeps ValueError as the single decode-failure contract that the docstring and the new boolList check already assume.

raise StepDetailsError(
action_id,
SessionActionLogKind.TASK_RUN,
str(e),
step_id=step_id,
task_id=task_id,
) from e

next_action = RunStepTaskAction(
id=action_id,
Expand Down
75 changes: 71 additions & 4 deletions src/deadline_worker_agent/sessions/job_entities/job_details.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,51 @@
from .validation import Field, validate_object


# The case-insensitive boolean vocabulary defined by the Open Job Description
# specification. Wire values are lowercased and membership-tested against these
# sets; keeping them as explicit sets (rather than a regex) makes the accepted
# tokens greppable and self-documenting.
_TRUE_STRINGS = frozenset({"true", "yes", "on", "1", "1.0"})
_FALSE_STRINGS = frozenset({"false", "no", "off", "0", "0.0"})

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

This hand-enumerated vocabulary is now a hard gate: any spelling the service emits that is not in these two sets fails the action (StepDetailsError for TASK_RUN, ValueError out of validate_entity_data for job details). Before this change the value was passed through verbatim and could not fail. So the sets have to match what the producer actually emits, exactly — and the contents suggest they may not.

The set mixes two different rules. true/false/yes/no/on/off is a token list, but 1, 1.0, 0, 0.0 look like the tail of a numeric rule. If the spec rule is numeric, then 1.00, 01, +1, 0.0000 are all equally valid spellings of the same booleans, and the tests here deliberately assert those are rejected (bool-1.00-rejected, bool-01-rejected, bool-00-rejected). If the rule is a token list, it is not obvious why 1.0/0.0 are members at all. One of the two readings is wrong, and if it is the strict one, workers will fail live tasks on parameter values the service considers valid.

Two things worth doing:

  1. Pin the accepted set against the actual producer rather than a prose reading of the spec — ideally openjd-models own boolean coercion if it exposes one (the repo already pins openjd-model >= 0.11.4, < 0.12, and this same decode feeds openjd.expr). Reusing it removes the divergence risk entirely and means the worker cannot drift when the spec adds a token. Re-deriving the vocabulary in the worker means every future spec addition is a worker-side hard-fail until a new agent ships — and agents in the field are not upgraded on the service is schedule.

  2. If the set must stay local, consider whether an unrecognized token should really be terminal. Given rollout timing, a permissive path (log a warning and fall back to the pre-change behavior for unknown spellings) fails softer than rejecting, since the downside of a wrong coercion is one mis-evaluated expression while the downside of rejection is a failed customer task.

Also worth noting: _TRUE_STRINGS / _FALSE_STRINGS overlap with two existing vocabularies in this same package — telemetry.py:30 (_TRUE_VALUES = {"true", "yes", "on", "1"}) and the config_file.py settings docs (0, off, f, false, n, no, 1, on, t, true, y, yes). Three different accepted sets for the same concept in one codebase is a drift hazard even if each is individually correct for its own input.



def _bool_from_api_response(value: str | bool) -> bool:
"""Coerces a wire-format boolean parameter value into a native Python bool.

The service is migrating boolean parameters from native JSON booleans to
string-typed booleans, but jobs created before the flip still carry native
booleans in persisted parameters that are returned verbatim, so both forms
may arrive during and after rollout. Coercing at the single wire-decode
choke point produces the native Python type that Open Job Description's
model documents expression-extension parameters to carry, and normalizes
both wire forms to one representation. It also validates the value, so an
out-of-vocabulary string fails fast here rather than reaching a consumer
that would render it verbatim.

A native bool is accepted and passed through unchanged. String values are
matched case-insensitively against the Open Job Description specification's
boolean vocabulary: "true", "yes", "on", "1", and "1.0" are True; "false",
"no", "off", "0", and "0.0" are False. Any other value -- including "maybe",
"", a string with surrounding whitespace, or a non-str, non-bool type such
as a native int 0 or 1 -- raises ValueError.
"""
# bool is a subclass of int, so match a native bool first and pass it
# through unchanged; the string membership tests below reject ints such as
# 0/1 because the vocabulary sets contain only str tokens.
if isinstance(value, bool):
return value
if isinstance(value, str):
lowered = value.lower()
if lowered in _TRUE_STRINGS:
return True
if lowered in _FALSE_STRINGS:
return False
raise ValueError(

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

This new ValueError is reachable from a code path where nothing catches it, so an unexpected boolean wire form will kill the whole session rather than failing one action.

parameters_from_api_response is called in two places:

  • JobDetails.from_boto (job_details.py:338), which is preceded by _validate_job_parameters — so a bad bool is rejected by the validator first, and that error is wrapped into a SessionActionError subclass by the caller. Fine.
  • SessionActionQueue.dequeue for TASK_RUN parameters (scheduler/session_queue.py:535). Task parameters are never passed through _validate_job_parameters, and that call sits outside the try/except (ValueError, RuntimeError) block above it (session_queue.py:516-533) which converts step_details failures into StepDetailsError.

So a bare ValueError propagates out of dequeue(). Session._start_action (sessions/session.py:654-658) only catches SessionActionError, so it escapes to Session._run → Session.run, which sets _stop_fail_message, sets the stop event and re-raises — tearing down the entire session and cancelling all remaining queued actions, instead of reporting a single FAILED action with a useful message.

Before this change, value["bool"] for a task parameter was passed through verbatim and could not raise, so this is a new failure mode — and it is precisely the one the rollout makes likely: if the service ever emits a bool form other than exactly true/false (or native bool) for a task parameter, workers hard-fail sessions.

Suggest moving the parameters_from_api_response call at session_queue.py:535 inside a try/except (ValueError, RuntimeError) that raises StepDetailsError(action_id, SessionActionLogKind.TASK_RUN, str(e), step_id=step_id, task_id=task_id), matching how step_details resolution failures are already handled a few lines above. That keeps the strict validation but degrades to a graceful per-action failure.

f"Expected a boolean parameter value of True, False, or one of the "
f"case-insensitive strings {sorted(_TRUE_STRINGS | _FALSE_STRINGS)} but got {value!r}"
)


def parameters_from_api_response(
params: dict[
str,
Expand Down Expand Up @@ -83,7 +128,9 @@ def parameters_from_api_response(
param_value = ParameterValue(type=ParameterValueType.CHUNK_INT, value=value["chunkInt"])
elif "bool" in value:
value = cast(BoolParameter, value)
param_value = ParameterValue(type=ParameterValueType.BOOL, value=value["bool"])
param_value = ParameterValue(
type=ParameterValueType.BOOL, value=_bool_from_api_response(value["bool"])

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

High-level question on the approach: after this change bool/boolList become the only parameter types in this function that are not passed through verbatim, and the coercion is what forces the new hard-fail behaviour. It is worth confirming the coercion is actually required before accepting that risk.

Look at what the surrounding branches do with the other scalar types whose wire form is a string:

  • int -> value["int"] passed through as a string (validator: isinstance(v, str))
  • float -> passed through as a string
  • chunkInt -> passed through as a string
  • rangeExpr -> passed through as a string
  • intList / floatList -> passed through as list[str] (validator: _is_str_list)

So the consumer downstream of ParameterValue already accepts un-parsed string forms for every other numeric type, including inside lists. If the service flip is normalizing bool to a string and boolList to list[str], that makes them consistent with int/intList rather than exceptional — which suggests the consumer may well accept them as-is too, and the docstring premise ("Open Job Description expression evaluation requires a native bool") may not hold, or may hold only for a path that native bool was already satisfying by accident before the flip.

This matters because the coercion is not free. It is what introduces the vocabulary gate, and therefore the new "unrecognized spelling fails the action" behaviour that did not exist before. If the string form is in fact accepted downstream, the lower-risk change is: widen the type annotations (which this PR already does), leave the value alone, and widen the validator to accept both forms — no coercion, no new failure mode, no worker-side copy of the spec vocabulary to keep in sync.

Two concrete things that would settle it:

  1. Point at where a native bool is required — openjd.models ParameterValue.value type, or the expression evaluator that consumes it. If ParameterValue.value is annotated str, then note that the pre-change code was passing a native bool into a str-typed field, and the coercion is preserving a type error rather than fixing one.
  2. Check the _v1/Rust path specifically: _to_rust_parameter_values (sessions/runtime/rust.py:122-152) hands value.value straight to the pyo3 TaskParameterValue. If that binding accepts the string form for BOOL the way it does for INT, the coercion buys nothing on the path that is presumably the migration target.

If the native bool genuinely is required, this all stands as-is and the only open item is the vocabulary itself (separate comment). If it is not, dropping the coercion removes the entire class of rollout risk.

)
elif "rangeExpr" in value:
value = cast(RangeExprParameter, value)
param_value = ParameterValue(
Expand All @@ -107,7 +154,18 @@ def parameters_from_api_response(
)
elif "boolList" in value:
value = cast(BoolListParameter, value)
param_value = ParameterValue(type=ParameterValueType.LIST_BOOL, value=value["boolList"])
bool_list = value["boolList"]
# Shape-check before iterating so a non-list (e.g. None) raises the
# same ValueError vocabulary the rest of this function uses rather
# than a TypeError that callers catching only ValueError would miss.
if not isinstance(bool_list, list):
raise ValueError(
f"Expected a boolList parameter value to be a list but got {bool_list!r}"
)
param_value = ParameterValue(
type=ParameterValueType.LIST_BOOL,
value=[_bool_from_api_response(item) for item in bool_list],
)
elif "intListList" in value:
value = cast(IntListListParameter, value)
param_value = ParameterValue(
Expand Down Expand Up @@ -499,19 +557,28 @@ def _validate_job_parameters(cls, job_parameters: dict[str, Any]) -> None:
def _is_str_list(value: Any) -> bool:
return isinstance(value, list) and all(isinstance(item, str) for item in value)

def _is_bool(value: Any) -> bool:
# Transitional: accept native JSON booleans (persisted before the
# flip) and the string form defined by the Open Job Description
# specification's case-insensitive boolean vocabulary. Matches what
# _bool_from_api_response coerces at wire-decode time.
if isinstance(value, bool):
return True
return isinstance(value, str) and value.lower() in _TRUE_STRINGS | _FALSE_STRINGS

value_shape_checks: dict[str, Callable[[Any], bool]] = {
"string": lambda v: isinstance(v, str),
"path": lambda v: isinstance(v, str),
"int": lambda v: isinstance(v, str),
"float": lambda v: isinstance(v, str),
"chunkInt": lambda v: isinstance(v, str),
"bool": lambda v: isinstance(v, bool),
"bool": _is_bool,
"rangeExpr": lambda v: isinstance(v, str),
"stringList": _is_str_list,
"pathList": _is_str_list,
"intList": _is_str_list,
"floatList": _is_str_list,
"boolList": lambda v: isinstance(v, list) and all(isinstance(item, bool) for item in v),
"boolList": lambda v: isinstance(v, list) and all(_is_bool(item) for item in v),
"intListList": lambda v: isinstance(v, list) and all(_is_str_list(item) for item in v),
}

Expand Down
201 changes: 201 additions & 0 deletions test/integ/sessions/test_differential_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@
from openjd.model import decode_environment_template, decode_job_template
from openjd.sessions import ActionStatus

from deadline_worker_agent.sessions.job_entities.job_details import (
parameters_from_api_response,
)
from deadline_worker_agent.sessions.runtime import (
SessionRuntime,
SessionRuntimeConfig,
Expand All @@ -21,8 +24,11 @@
)

if TYPE_CHECKING:
from openjd.model import ParameterValue
from openjd.sessions import EnvironmentModel, StepScriptModel

from deadline_worker_agent.api_models import BoolParameter

# These tests run the SAME scenarios through BOTH the Python (v0) and Rust (v1)
# SessionRuntime adapters against real openjd sessions executing real local
# subprocesses. They require no AWS/farm resources: every scenario is driven by
Expand Down Expand Up @@ -207,6 +213,92 @@ def _wait_for_terminal(runtime: SessionRuntime, timeout: float = 30.0) -> Option
return _state_name(runtime)


def _resolve_job_param_reference(
runtime_kind: SessionRuntimeKind,
*,
job_parameter_values: dict[str, ParameterValue],
param_name: str,
param_type: str,
root: Path,
) -> str:
"""Resolve a single ``{{Param.<param_name>}}`` reference through a real session.

Builds a runtime of the requested kind with the EXPR extension enabled --
the boolean parameter types (BOOL, LIST[BOOL]) are EXPR-extension types and
the decoder rejects them without it -- then enters an environment whose
onEnter writes the resolved reference to a file and returns the captured
text. The reference is spliced into the ``-c`` body as a single-quoted
Python string literal, so a LIST[BOOL] value is observed as one rendered
string (e.g. ``[true, false]``) rather than expanded into separate process
arguments the way a bare list reference in ``args`` would be.

The environment template both declares ``extensions: ["EXPR"]`` and is
decoded with ``supported_extensions=["EXPR"]``; both are required -- the
template field opts the template in, and the decode argument is the
implementation allowlist the field is intersected against.
"""
output = root / "resolved_param.txt"
reference = f"{{{{Param.{param_name}}}}}"
recorder = _StatusRecorder()
config = SessionRuntimeConfig(
session_id=f"session-boolparam-{runtime_kind.name.lower()}",
job_parameter_values=job_parameter_values,
path_mapping_rules=None,
retain_working_dir=False,
user=None,
action_callback=recorder,
os_env_vars=None,
session_root_directory=root,
supported_extensions=("EXPR",),
)
runtime = create_session_runtime(runtime_kind, config)
try:
env_template = decode_environment_template(
template={
"specificationVersion": "environment-2023-09",
"extensions": ["EXPR"],
"environment": {
"name": "BoolParamEnv",
"script": {
"actions": {
"onEnter": {
"command": sys.executable,
"args": [
"-c",
f"import sys; open(sys.argv[1], 'w').write('{reference}')",
str(output),
],
}
}
},
},
"parameterDefinitions": [{"name": param_name, "type": param_type}],
},
supported_extensions=["EXPR"],
)

prev = len(recorder.statuses)
identifier = runtime.enter_environment(environment=env_template.environment)
assert _wait_for_new_action(recorder, prev)
terminal = _wait_for_terminal(runtime)
assert terminal == "SUCCESS", (
f"{runtime_kind.name}: expected SUCCESS resolving Param.{param_name} "
f"but got {terminal}; statuses={recorder.state_names}"
)
resolved = output.read_text()

prev = len(recorder.statuses)
runtime.exit_environment(identifier=identifier)
assert _wait_for_new_action(recorder, prev)
assert _wait_for_terminal(runtime) == "SUCCESS"
return resolved
finally:
try:
runtime.cleanup()
except Exception:
pass


class TestDifferentialSessionRuntime:
"""Differential behavior tests: Python (v0) vs Rust (v1) adapters.

Expand Down Expand Up @@ -452,3 +544,112 @@ def test_enter_environment_when_job_has_typed_params_resolves_param_reference(
runtime.cleanup()
except Exception:
pass

# ------------------------------------------------------------------
# Boolean wire-format parameter resolution.
#
# The worker decodes boolean task/job parameters at a single wire-decode
# choke point (``parameters_from_api_response``), coercing the two wire
# forms the service may send -- the native JSON boolean ``{"bool": true}``
# (sent today) and the string ``{"bool": "true"}`` (sent after the model
# change) -- into a native Python bool. These tests exercise that decode
# path end to end: a WIRE-FORMAT dict goes through
# ``parameters_from_api_response`` and the decoded value is fed to a real
# session that references it via ``{{Param.X}}``, so the resolved value is
# observed in output rather than merely type-checked at the decode boundary.
#
# Both wire forms, and every non-canonical-but-legal token, must resolve to
# the canonical lowercase ``true``/``false`` OpenJD renders for a bool, and
# must do so identically on the Python (v0) and Rust (v1) runtimes. The
# expected strings below are spec-derived literals, not values recomputed
# from the decoder, and each parametrized runtime asserts against the same
# literal -- so a runtime that rendered a bool differently (e.g. ``True`` or
# ``1``) would fail its own case rather than be averaged away.

@pytest.mark.timeout(60)
@pytest.mark.parametrize("runtime_kind", _RUNTIME_KINDS)
@pytest.mark.parametrize(
"wire_value, expected",
[
pytest.param({"bool": True}, "true", id="native-true"),
pytest.param({"bool": "true"}, "true", id="string-true"),
pytest.param({"bool": "yes"}, "true", id="string-yes"),
pytest.param({"bool": "1"}, "true", id="string-1"),
pytest.param({"bool": False}, "false", id="native-false"),
pytest.param({"bool": "0"}, "false", id="string-0"),
],
)
def test_bool_job_param_wire_form_resolves_to_canonical_string(
self,
runtime_kind: SessionRuntimeKind,
wire_value: BoolParameter,
expected: str,
tmp_path: Path,
) -> None:
root = tmp_path / f"boolparam-{runtime_kind.name.lower()}"
root.mkdir(parents=True, exist_ok=True)

# The API-shaped dict is decoded through the real wire path -- the same
# code the worker runs on a BatchGetJobEntity response -- not by
# hand-building a ParameterValue.
job_parameter_values = parameters_from_api_response({"MyBool": wire_value})

resolved = _resolve_job_param_reference(
runtime_kind,
job_parameter_values=job_parameter_values,
param_name="MyBool",
param_type="BOOL",
root=root,
)

assert resolved == expected

@pytest.mark.timeout(60)
@pytest.mark.parametrize("runtime_kind", _RUNTIME_KINDS)
def test_bool_list_job_param_wire_form_resolves_to_canonical_strings(
self, runtime_kind: SessionRuntimeKind, tmp_path: Path
) -> None:
root = tmp_path / f"boollist-{runtime_kind.name.lower()}"
root.mkdir(parents=True, exist_ok=True)

# A LIST[BOOL] mixing the native form with non-canonical legal tokens;
# every element must decode and resolve to its canonical rendering.
job_parameter_values = parameters_from_api_response(
{"MyBools": {"boolList": [True, "false", "yes", "0"]}}
)

resolved = _resolve_job_param_reference(
runtime_kind,
job_parameter_values=job_parameter_values,
param_name="MyBools",
param_type="LIST[BOOL]",
root=root,
)

assert resolved == "[true, false, true, false]"

@pytest.mark.timeout(90)
def test_bool_param_resolution_is_identical_across_runtimes(self, tmp_path: Path) -> None:
"""The same wire-form bool must resolve to the same concrete text on both
runtimes. This compares them directly rather than relying on each
asserting a shared literal, so a silent divergence -- one runtime
rendering ``True`` or ``1`` while the other renders ``true`` -- fails
here. A non-canonical token (``"yes"``) is used so the assertion also
pins that the decoder's vocabulary is applied consistently on both."""
wire_value: BoolParameter = {"bool": "yes"}
outputs: dict[str, str] = {}
for runtime_kind in (SessionRuntimeKind.PYTHON, SessionRuntimeKind.RUST):
root = tmp_path / f"cross-{runtime_kind.name.lower()}"
root.mkdir(parents=True, exist_ok=True)
job_parameter_values = parameters_from_api_response({"MyBool": wire_value})
outputs[runtime_kind.name] = _resolve_job_param_reference(
runtime_kind,
job_parameter_values=job_parameter_values,
param_name="MyBool",
param_type="BOOL",
root=root,
)

assert outputs["PYTHON"] == outputs["RUST"]
# And the shared value is the canonical rendering, not merely equal-but-wrong.
assert outputs["PYTHON"] == "true"
Loading
Loading