Files
DynamisLab/tests/test_drl_pinball_v5_acquire.py
T
Frank14fandCursor 61e82ec90a feat(eval): publish cycle-mean wake acquisition
Add deterministic phase-filtered V5 and Legacy acquisition with complete-cycle mean fields, then evaluate controlled wakes against target and zero baselines offline.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-08 15:50:49 +08:00

566 lines
26 KiB
Python

from __future__ import annotations
import subprocess
import sys
from pathlib import Path
import numpy as np
import pytest
from drl_pinball.eval import acquire_v5
class FakeModel:
def __init__(self):
self.calls = 0
def predict(self, obs, deterministic):
assert deterministic is True
self.calls += 1
return np.array([[0.1, -0.2, 0.3]], dtype=np.float32), None
class FakeVecEnv:
def __init__(self, raw):
self.raw = raw
self.steps = 0
self.events = []
def reset(self):
self.events.append("reset")
return np.zeros((1, 12), dtype=np.float32)
def step(self, action):
self.steps += 1
self.raw.control_step = self.steps
self.raw.sim.stepper.step_count = self.steps * 800
self.raw.smoother._state = self.raw._action_to_omega(action) * 0.5
self.events.append(("step", self.steps))
info = {"sim": self.steps / 1000, "r_cd": 1, "r_cl": 2, "r_sim": 3, "floor_pen": 4}
return np.zeros((1, 12)), np.array([5.0]), np.array([False]), [info]
class FakeRaw:
def __init__(self):
self.control_step = 0
self._cal = {"U0": 0.01, "grid": {"nx": 2000, "ny": 600},
"ACTION_BIAS": [0.0, 0.0, 0.0], "ACTION_SCALE": 12.0}
self.smoother = type("Smoother", (), {"_state": np.zeros(3)})()
self.sim = type("Sim", (), {})()
self.sim.stepper = type("Stepper", (), {"step_count": 0})()
def _action_to_omega(self, action):
return np.asarray(action).reshape(3) * 2
def _read_obs(self):
return np.arange(14, dtype=np.float32) + self.control_step
def test_controlled_schedule_is_reset_750_warmup_then_250_post_step_fields(tmp_path):
raw, model = FakeRaw(), FakeModel()
vec = FakeVecEnv(raw)
captures = []
def capture(env):
captures.append((env.control_step, len(vec.events)))
value = np.full((2, 3), env.control_step, dtype=np.float32)
return {"rho": value, "ux": value, "uy": value}
rows, buffer = acquire_v5._collect_controlled(model, vec, raw, tmp_path, capture)
assert vec.events[0] == "reset"
assert vec.steps == model.calls == 1000
assert len(rows) == len(captures) == 250
assert buffer["ux"].shape == buffer["uy"].shape == (250, 2, 3)
assert captures[0][0] == rows[0]["control_index"] == 751
assert captures[-1][0] == rows[-1]["control_index"] == 1000
assert all(event_count == step + 1 for step, event_count in captures)
assert rows[0]["native_reward_dtw"] == pytest.approx(0.751)
assert np.allclose(rows[0]["commanded_target_omega"], [0.2, -0.4, 0.6])
assert np.allclose(rows[0]["effective_smoothed_omega"], [0.1, -0.2, 0.3])
assert np.allclose(buffer["ux"][0], 751) and np.allclose(buffer["uy"][-1], 1000)
assert not list(tmp_path.glob("boundary_*.npz"))
def test_field_capture_runs_inside_env_cuda_context_and_validates_shape():
events = []
raw = type("Raw", (), {})()
raw.sim = type("Sim", (), {})()
raw.sim.lbm_cfg = type("Cfg", (), {"nx": 3, "ny": 2})()
def macro():
events.append("macro")
value = np.ones((2, 3), dtype=np.float32)
return {"rho": value, "ux": value, "uy": value}
raw.sim.get_macroscopic = macro
raw._gpu_block = lambda fn: (events.append("push"), fn(), events.append("pop"))
result = acquire_v5._capture_fields(raw)
assert events == ["push", "macro", "pop"]
assert result["ux"].shape == (2, 3)
assert set(result) == {"rho", "ux", "uy"}
def test_zero_uses_full_vec_step_schedule_and_zero_action(tmp_path):
raw = FakeRaw()
vec = FakeVecEnv(raw)
rows, buffer = acquire_v5._collect_zero(vec, raw, tmp_path, lambda env: {
name: np.ones((2, 3), dtype=np.float32) for name in ("rho", "ux", "uy")
})
assert vec.events[0] == "reset"
assert vec.steps == 1000 and len(rows) == 250
assert buffer["ux"].shape == (250, 2, 3)
assert np.array_equal(rows[0]["action_normalized"], np.zeros(3, dtype=np.float32))
assert rows[0]["native_reward_dtw"] == pytest.approx(0.751)
assert rows[-1]["control_index"] == 1000
assert set(rows[0]) == set(acquire_v5._target_boundary(type("Target", (), {
"sensor_ids": (0, 1, 2), "calibration": {"U0": 0.01, "grid": {"nx": 2000}},
"sim": type("Sim", (), {"stepper": type("Stepper", (), {"step_count": 800})(),
"read_sensor": lambda self, sid, normalize: (0.0, 0.0)})()
})(), 1))
def test_target_geometry_schedule_order_and_nan_contract(tmp_path):
class Sim:
def __init__(self):
self.added, self.runs, self.closed = [], [], False
self._objects = []
self.bodies = type("Bodies", (), {
"get": lambda owner, index: self._objects[index],
"count": property(lambda owner: len(self._objects)),
})()
self.stepper = type("Stepper", (), {"step_count": 0})()
self.lbm_cfg = type("Cfg", (), {"nx": 3, "ny": 2})()
context = type("Context", (), {"push": lambda self: None, "pop": lambda self: None})()
self.ctx = type("Cuda", (), {"_ctx": context})()
def add_body(self, kind, **kwargs):
self.added.append((kind, kwargs))
body_id = len(self.added) - 1
self._objects.append(type("Body", (), {
"obj_id": body_id, "_is_sensor": kind == "sensor",
})())
return body_id
def initialize(self):
self.initialized = True
def run(self, steps, **kwargs):
self.runs.append((steps, kwargs))
self.stepper.step_count += steps
def read_sensor(self, sensor_id, normalize=True):
assert normalize is True
return np.array([sensor_id + 0.1, sensor_id + 0.2])
def get_macroscopic(self):
value = np.ones((2, 3), dtype=np.float32)
return {"rho": value, "ux": value, "uy": value}
def close(self):
self.closed = True
sim = Sim()
bundle = {"calibration": {"grid": {"nx": 3, "ny": 2}, "dist_radius": 1.25,
"L0": 20.0, "U0": 0.01},
"config_path": Path("config.json")}
case = type("Case", (), {"scene_type": "karman", "target_diam": None})()
spinups = []
runtime = acquire_v5._create_target_runtime(
case, bundle, 2, simulation_factory=lambda **_: sim,
spinup_runner=lambda target_sim, steps: spinups.append((target_sim, steps)),
)
assert sim.added == [
("circle", {"center": (600.0, 0.5, 0.0), "radius": 25.0}),
("sensor", {"center": (1200.0, 40.5, 0.0), "radius": 5.0}),
("sensor", {"center": (1200.0, 0.5, 0.0), "radius": 5.0}),
("sensor", {"center": (1200.0, -39.5, 0.0), "radius": 5.0}),
]
assert spinups == [(sim, 1200)]
rows, buffer = acquire_v5._collect_target(runtime, tmp_path, 800)
assert sim.runs == [(800, {"zero_obs": True, "sync_obs": True})] * 1000
assert len(rows) == 250
assert buffer["ux"].shape == buffer["uy"].shape == (250, 2, 3)
assert np.allclose(buffer["ux"][0], 1.0)
assert not list(tmp_path.glob("boundary_*.npz"))
assert np.allclose(rows[0]["sensors"], [1.1, 1.2, 2.1, 2.2, 3.1, 3.2])
for name in ("forces", "action_normalized", "commanded_target_omega",
"effective_smoothed_omega"):
assert np.all(np.isnan(rows[0][name]))
for name in ("reward_raw", "cd", "cl", "r_cd", "r_cl", "r_sim", "floor_pen",
"native_reward_dtw"):
assert np.isnan(rows[0][name])
runtime.close()
assert sim.closed
def test_finalize_converts_sensors_only_for_dtw(tmp_path, monkeypatch):
captured = {}
n = 150
times = np.arange(n, dtype=float)
sensors = np.column_stack([np.sin(2 * np.pi * times / 30 + i) for i in range(6)])
rows = []
for i in range(n):
rows.append({"physical_time": float(i), "lattice_step": i * 800,
"control_index": i + 1, "sensors": sensors[i], "forces": np.ones(6),
"action_normalized": np.zeros(3), "commanded_target_omega": np.zeros(3),
"effective_smoothed_omega": np.zeros(3), "reward_raw": 1.0,
"cd": 1.0, "cl": 1.0, "r_cd": 1.0, "r_cl": 1.0, "r_sim": 1.0,
"floor_pen": 0.0, "native_reward_dtw": 1.0})
scratch_root = tmp_path / "scratch"
scratch = scratch_root / "candidate"
scratch.mkdir(parents=True)
fields = {
"ux": np.ones((n, 2, 3), dtype=np.float32),
"uy": np.ones((n, 2, 3), dtype=np.float32),
}
identity = tmp_path / "identity"
identity.write_bytes(b"read-only")
bundle = {"target_states": sensors * 7.0, "model_path": identity,
"vecnormalize_path": identity, "config_path": identity}
original = acquire_v5.dual_cycle_dtw
def observe(target, state, native, **kwargs):
captured["state"] = state.copy()
captured["lag_channel"] = kwargs["lag_channel"]
return original(target, state, native, **kwargs)
monkeypatch.setattr(acquire_v5, "dual_cycle_dtw", observe)
monkeypatch.setattr(acquire_v5.infer_train, "_file_identity", lambda path: {"path": str(path)})
monkeypatch.setattr(acquire_v5.infer_train, "_bundle_metadata", lambda bundle: {})
case = type("Case", (), {"case_id": "kar_re100", "si": 800})()
acquire_v5._finalize(tmp_path, scratch, rows, fields, bundle,
{"resolved_output_root": tmp_path}, 7.0, "zero", [],
case=case, seed=45, cycle_length=30)
assert np.allclose(captured["state"], sensors * 7.0)
assert captured["lag_channel"] == 3
import json
assert json.loads((tmp_path / "dtw_summary.json").read_text())["lag_channel"] == 3
assert json.loads((tmp_path / "metadata.json").read_text())["dtw_lag_channel"] == 3
with np.load(tmp_path / "timeseries.npz", allow_pickle=False) as saved:
assert np.allclose(saved["sensors"], sensors)
assert identity.read_bytes() == b"read-only"
def test_collection_failure_cleans_only_transaction_scratch(tmp_path):
scratch_root = tmp_path / "scratch"
scratch_root.mkdir()
scratch = acquire_v5.create_scratch(scratch_root)
sibling = tmp_path / "immutable-model.zip"
sibling.write_bytes(b"model")
(scratch / "partial.npz").write_bytes(b"partial")
acquire_v5.cleanup_scratch(scratch, root=scratch_root)
assert not scratch.exists()
assert sibling.read_bytes() == b"model"
def test_acquire_finalize_failure_leaves_no_partial_role(tmp_path, monkeypatch):
final_role = tmp_path / "v5" / "karman_re100" / "controlled"
sentinel = tmp_path / "immutable-model.zip"
sentinel.write_bytes(b"model")
storage = {"resolved_output_root": tmp_path, "device": tmp_path.stat().st_dev}
bundle = {"model_path": sentinel, "vecnormalize_path": sentinel}
monkeypatch.setattr(acquire_v5, "get_case", lambda _: type(
"Case", (), {"case_id": "kar_re100", "scene_type": "karman", "si": 800,
"seeds": (45,)})())
monkeypatch.setattr(acquire_v5.infer_train, "_resolve_seed_artifacts", lambda *_: bundle)
monkeypatch.setattr(acquire_v5, "_validate_acquisition_bundle", lambda *_: 30)
monkeypatch.setattr(acquire_v5, "_validate_shared_role_identity", lambda *_: None)
monkeypatch.setattr(acquire_v5, "_collect_controlled", lambda *_ , **__: ([], []))
class Env:
def close(self):
pass
runtime = lambda *_: (Env(), object(), object(), type("Raw", (), {})())
def fail(staging, *_, **__):
(staging / "timeseries.npz").write_bytes(b"partial")
raise RuntimeError("injected finalize failure")
with pytest.raises(RuntimeError, match="injected finalize failure"):
acquire_v5.acquire_controlled(
output_root=tmp_path, overwrite=True,
storage_validator=lambda **_: storage, runtime_factory=runtime, finalizer=fail,
)
assert not final_role.exists()
case_dir = final_role.parent
assert not case_dir.exists() or list(case_dir.iterdir()) == []
assert sentinel.read_bytes() == b"model"
def test_target_full_finalize_exact_products_and_unavailable_metadata(tmp_path, monkeypatch):
n = 150
times = np.arange(n, dtype=float)
sensors = np.column_stack([np.sin(2 * np.pi * times / 30 + i) for i in range(6)])
rows = []
for i in range(n):
nan3, nan6 = np.full(3, np.nan), np.full(6, np.nan)
rows.append({"physical_time": float(i), "lattice_step": i * 800,
"control_index": i + 1, "sensors": sensors[i], "forces": nan6,
"action_normalized": nan3, "commanded_target_omega": nan3,
"effective_smoothed_omega": nan3, "reward_raw": np.nan,
"cd": np.nan, "cl": np.nan, "r_cd": np.nan, "r_cl": np.nan,
"r_sim": np.nan, "floor_pen": np.nan, "native_reward_dtw": np.nan})
role_dir = tmp_path / "role"
role_dir.mkdir()
scratch = role_dir / "scratch" / "candidate"
scratch.mkdir(parents=True)
fields = {
"ux": np.ones((n, 2, 3), dtype=np.float32),
"uy": np.ones((n, 2, 3), dtype=np.float32),
}
identity = tmp_path / "identity"
identity.write_bytes(b"read-only")
unavailable = ["forces", "action_normalized", "reward_raw", "native_reward_dtw"]
monkeypatch.setattr(acquire_v5.infer_train, "_file_identity", lambda path: {"path": str(path)})
monkeypatch.setattr(acquire_v5.infer_train, "_bundle_metadata", lambda bundle: {})
acquire_v5._finalize(
role_dir, scratch, rows, fields,
{"target_states": sensors, "model_path": identity, "vecnormalize_path": identity,
"config_path": identity},
{"resolved_output_root": tmp_path}, 1.0, "target", unavailable,
case=type("Case", (), {"case_id": "kar_re100", "si": 800})(),
seed=45, cycle_length=30,
)
monkeypatch.setattr(acquire_v5, "COLLECT_BOUNDARIES", n)
acquire_v5._validate_staged_role(role_dir)
expected = {"timeseries.npz", "timeseries.csv", "phase_cycle.npz", "phase_cycle.csv",
"phase_fields.npz", "dtw_summary.json", "metadata.json", "identity"}
expected.remove("identity")
assert {path.name for path in role_dir.iterdir()} == expected
import json
metadata = json.loads((role_dir / "metadata.json").read_text())
summary = json.loads((role_dir / "dtw_summary.json").read_text())
assert metadata["role"] == "target" and metadata["seed"] is None
assert metadata["unavailable_fields"] == unavailable
assert metadata["candidate_field_storage"].startswith("single-role in-memory")
assert metadata["phase_smoothing_kernel"] == [0.25, 0.5, 0.25]
assert metadata["mean_field_count"] > 0
assert summary["native_mean"] is None
with np.load(role_dir / "timeseries.npz", allow_pickle=False) as saved:
assert np.all(np.isnan(saved["native_reward_dtw"]))
assert np.allclose(saved["sensors"], sensors)
with np.load(role_dir / "phase_cycle.npz", allow_pickle=False) as saved:
assert "sensors_pooled" in saved.files and "reward_raw_mean" in saved.files
assert len(saved["sensors_pooled"]) > 0
with np.load(role_dir / "phase_fields.npz", allow_pickle=False) as saved:
assert set(saved.files) == acquire_v5.PHASE_FIELD_KEYS
assert saved["mean_ux"].shape == saved["ux"].shape[1:]
assert saved["mean_uy"].shape == saved["uy"].shape[1:]
def test_cli_enables_all_roles_without_replay_rejection():
source = Path(acquire_v5.__file__).read_text()
assert 'parser.add_argument("--case", choices=CASE_IDS' in source
assert 'parser.add_argument("--seed", type=int)' in source
assert 'parser.add_argument("--role", choices=ROLES' in source
assert "replay is not implemented" not in source
assert "acquire_role(args.role" in source
def test_target_unavailable_summary_is_strict_json(tmp_path):
path = tmp_path / "summary.json"
acquire_v5._atomic_json(path, {"native_mean": None, "unavailable_fields": ["reward_raw"]})
text = path.read_text()
assert "NaN" not in text and '"native_mean": null' in text
def test_scene_aware_raw_sample_layouts():
karman = type("Raw", (), {"_read_obs": lambda self: np.arange(14, dtype=np.float32)})()
illusion = type("Raw", (), {"_read_obs": lambda self: np.arange(12, dtype=np.float32)})()
assert np.array_equal(acquire_v5._raw_sample(karman, "karman"), np.arange(2, 14))
assert np.array_equal(acquire_v5._raw_sample(illusion, "illusion"), np.arange(12))
with pytest.raises(ValueError, match="6-sensor/6-force"):
acquire_v5._raw_sample(illusion, "karman")
def test_illusion_target_geometry_and_case_si(tmp_path):
class Sim:
def __init__(self):
self.added, self.runs = [], []
self._objects = []
self.bodies = type("Bodies", (), {
"get": lambda owner, index: self._objects[index],
"count": property(lambda owner: len(self._objects)),
})()
self.stepper = type("Stepper", (), {"step_count": 0})()
self.lbm_cfg = type("Cfg", (), {"nx": 3, "ny": 2})()
def add_body(self, kind, **kwargs):
self.added.append((kind, kwargs))
body_id = len(self.added) - 1
self._objects.append(type("Body", (), {
"obj_id": body_id, "_is_sensor": kind == "sensor",
})())
return body_id
def initialize(self): pass
def run(self, steps, **kwargs):
self.runs.append(steps); self.stepper.step_count += steps
def read_sensor(self, sensor_id, normalize=True): return (0.1, 0.2)
def get_macroscopic(self):
value = np.ones((2, 3), dtype=np.float32)
return {name: value for name in ("rho", "ux", "uy")}
def close(self): pass
sim = Sim()
context = type("Context", (), {"push": lambda self: None, "pop": lambda self: None})()
sim.ctx = type("Cuda", (), {"_ctx": context})()
case = type("Case", (), {"scene_type": "illusion", "target_diam": 1.5})()
bundle = {"calibration": {"grid": {"nx": 3, "ny": 2}, "L0": 20.0, "U0": 0.01},
"config_path": Path("config.json")}
spinups = []
runtime = acquire_v5._create_target_runtime(
case, bundle, 0, lambda **_: sim,
spinup_runner=lambda target_sim, steps: spinups.append((target_sim, steps)),
)
assert sim.added[0] == ("circle", {"center": (400.0, 0.5, 0.0), "radius": 30.0})
assert [item[1]["center"][0] for item in sim.added[1:]] == [600.0] * 3
assert spinups == [(sim, 1200)]
acquire_v5._collect_target(runtime, tmp_path, 1200)
assert sim.runs == [1200] * 1000
def test_bundle_validation_covers_registry_and_fails_before_storage(tmp_path, monkeypatch):
target = np.zeros((150, 6), dtype=np.float32)
phase = 2 * np.pi * np.arange(150) / 30
target[:, 3] = np.sin(phase)
config = tmp_path / "config.json"
calibration = tmp_path / "calibration.json"
config.write_text('{"grid":{"nx":2000,"ny":600},"physics":{"velocity":0.01}}')
calibration.write_text('{"SI":800}')
case = type("Case", (), {
"case_id": "kar_re100", "scene_type": "karman", "si": 800,
"seeds": (45,), "target_diam": None, "config_path": config,
})()
bundle = {"seed": "45", "config_path": config, "calibration_path": calibration,
"calibration": {"SI": 800, "U0": 0.01, "grid": {"nx": 2000, "ny": 600}},
"target_states": target}
assert set(acquire_v5.CYCLE_WINDOWS) == set(acquire_v5.CASE_IDS)
assert acquire_v5._validate_acquisition_bundle(case, 45, bundle) == 30
calibration.write_text('{"SI":500}')
with pytest.raises(ValueError, match="SI"):
acquire_v5._validate_acquisition_bundle(case, 45, bundle)
def test_output_paths_seed_qualify_controlled_only(tmp_path, monkeypatch):
case = type("Case", (), {"case_id": "kar_re100", "scene_type": "karman",
"si": 800, "seeds": (45,), "target_diam": None})()
bundle = {"seed": "45"}
storage = {"resolved_output_root": tmp_path, "device": tmp_path.stat().st_dev}
monkeypatch.setattr(acquire_v5, "get_case", lambda _: case)
monkeypatch.setattr(acquire_v5.infer_train, "_resolve_seed_artifacts", lambda *_: bundle)
monkeypatch.setattr(acquire_v5, "_validate_acquisition_bundle", lambda *_: 30)
monkeypatch.setattr(acquire_v5, "_validate_staged_role", lambda *_: None)
monkeypatch.setattr(acquire_v5, "publish_role_output", lambda prepared: prepared["final_role_dir"])
class Env:
def close(self): pass
monkeypatch.setattr(acquire_v5, "_collect_controlled", lambda *_, **__: ([], []))
def finalize(role_dir, scratch, *args, **kwargs):
acquire_v5.cleanup_scratch(scratch, root=role_dir / "scratch")
scratch.parent.rmdir()
result = acquire_v5.acquire_role(
"controlled", case_id="kar_re100", seed=45, output_root=tmp_path,
storage_validator=lambda **_: storage,
runtime_factory=lambda *_: (Env(), object(), object(), object()), finalizer=finalize,
)
assert result == tmp_path / "v5/kar_re100_seed45/controlled"
def test_shared_roles_require_seed_invariant_physical_identity(monkeypatch):
case = type("Case", (), {"case_id": "kar_re100", "seeds": (41, 42)})()
target = np.ones((150, 6), dtype=np.float32)
base = {"seed": "41", "calibration": {"SI": 800, "config_path": "old"},
"target_states": target, "config_path": Path("config.json")}
other = {"seed": "42", "calibration": {"SI": 800, "config_path": "new"},
"target_states": target.copy(), "config_path": Path("config.json")}
monkeypatch.setattr(acquire_v5, "_validate_acquisition_bundle", lambda *_: 30)
monkeypatch.setattr(acquire_v5.infer_train, "_resolve_seed_artifacts", lambda *_: other)
acquire_v5._validate_shared_role_identity(case, 41, base)
other["target_states"] = target + np.float32(1e-3)
with pytest.raises(ValueError, match="different physical target"):
acquire_v5._validate_shared_role_identity(case, 41, base)
def test_target_runtime_fails_closed_on_body_id_order_and_count():
class Bodies:
def __init__(self, sim): self.sim = sim
@property
def count(self): return len(self.sim.objects)
def get(self, index): return self.sim.objects[index]
class BadSim:
def __init__(self):
self.objects = []
self.bodies = Bodies(self)
def add_body(self, kind, **kwargs):
body_id = len(self.objects) + 1
self.objects.append(type("Body", (), {
"obj_id": body_id, "_is_sensor": kind == "sensor",
})())
return body_id
def initialize(self): raise AssertionError("must fail before initialize")
case = type("Case", (), {"scene_type": "karman", "target_diam": None})()
bundle = {"calibration": {"grid": {"nx": 2000, "ny": 600}, "U0": 0.01,
"L0": 20.0}, "config_path": Path("config.json")}
with pytest.raises(ValueError, match="body order"):
acquire_v5._create_target_runtime(case, bundle, 0, lambda **_: BadSim())
def test_physical_zero_counterbias_shape_and_rollout(tmp_path):
class BiasedRaw(FakeRaw):
def __init__(self):
super().__init__()
self._cal.update(ACTION_BIAS=[1.5, -3.0, 0.75], ACTION_SCALE=6.0)
def _action_to_omega(self, action):
action = np.asarray(action, dtype=np.float32).reshape(3)
return action * self._cal["ACTION_SCALE"] + np.asarray(
self._cal["ACTION_BIAS"], dtype=np.float32)
raw = BiasedRaw()
vec = FakeVecEnv(raw)
actions = []
original_step = vec.step
def step(action):
actions.append(np.asarray(action).copy())
return original_step(action)
vec.step = step
acquire_v5._collect_zero(vec, raw, tmp_path, lambda env: {
name: np.ones((2, 3), dtype=np.float32) for name in ("rho", "ux", "uy")
})
expected = np.array([[-0.25, 0.5, -0.125]], dtype=np.float32)
assert actions and all(action.shape == (1, 3) for action in actions)
assert all(np.array_equal(action, expected) for action in actions)
def test_acquire_role_propagates_case_seed_si_and_scene(tmp_path, monkeypatch):
case = type("Case", (), {"case_id": "ill_1L", "scene_type": "illusion",
"si": 1200, "seeds": (43,), "target_diam": 1.0})()
bundle = {"seed": "43"}
storage = {"resolved_output_root": tmp_path, "device": tmp_path.stat().st_dev}
observed = {}
monkeypatch.setattr(acquire_v5, "get_case", lambda case_id: case)
monkeypatch.setattr(acquire_v5.infer_train, "_resolve_seed_artifacts",
lambda selected_case, seed: bundle)
monkeypatch.setattr(acquire_v5, "_validate_acquisition_bundle", lambda *args: 19)
monkeypatch.setattr(acquire_v5, "_validate_staged_role", lambda *_: None)
monkeypatch.setattr(acquire_v5, "publish_role_output", lambda prepared: prepared["final_role_dir"])
class Env:
def close(self): pass
def collect(model, vec, raw, scratch, **kwargs):
observed["scene_type"] = kwargs["scene_type"]
return [], []
monkeypatch.setattr(acquire_v5, "_collect_controlled", collect)
def finalize(role_dir, scratch, *args, **kwargs):
observed.update(case=kwargs["case"], seed=kwargs["seed"], cycle=kwargs["cycle_length"])
acquire_v5.cleanup_scratch(scratch, root=role_dir / "scratch")
scratch.parent.rmdir()
result = acquire_v5.acquire_role(
"controlled", case_id="ill_1L", seed=43, output_root=tmp_path,
storage_validator=lambda **_: storage,
runtime_factory=lambda selected_case, selected_bundle, device: (
Env(), object(), object(), type("Raw", (), {"_dtw_sensor_factor": 78.0})()),
finalizer=finalize,
)
assert observed == {"scene_type": "illusion", "case": case, "seed": 43, "cycle": 19}
assert result == tmp_path / "v5/ill_1L_seed43/controlled"