diff --git a/tools/hwlab-cli/caserun-runtime-observation.test.ts b/tools/hwlab-cli/caserun-runtime-observation.test.ts index f0264099..5b954c58 100644 --- a/tools/hwlab-cli/caserun-runtime-observation.test.ts +++ b/tools/hwlab-cli/caserun-runtime-observation.test.ts @@ -1,4 +1,5 @@ import assert from "node:assert/strict"; +import { spawnSync } from "node:child_process"; import { mkdtemp, rm, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; @@ -63,10 +64,27 @@ test("HWPOD node owns ioProbe command output normalization before CaseRun projec document.spec.ioProbe.endpoints.main41.command = `${process.execPath} ${fakeBoardComm}`; const plan = compileHwpodNodeOpsPlan({ document, intent: "io.probe.read", args: { probeId: "main41.ai0.current", count: 3 } }); const nodeResult = await executeHwpodNodeOpsPlan(plan, { now: () => "2026-07-16T00:00:00.000Z" }); + const pythonScript = ` +import importlib.util, json, logging, sys +spec = importlib.util.spec_from_file_location("hwlab_node", "tools/hwlab-node.py") +module = importlib.util.module_from_spec(spec) +spec.loader.exec_module(module) +plan = json.loads(sys.argv[1]) +executor = module.NodeOpsExecutor({"nodeId": plan["nodeId"]}, logging.getLogger("test")) +print(json.dumps(executor.execute(plan))) +`; + const python = spawnSync("python3", ["-c", pythonScript, JSON.stringify(plan)], { cwd: process.cwd(), encoding: "utf8" }); + assert.equal(python.status, 0, python.stderr); + const pythonResult = JSON.parse(python.stdout); assert.equal(nodeResult.ok, true); assert.equal(nodeResult.results[0].output.observation.probeId, "main41.ai0.current"); assert.deepEqual(nodeResult.results[0].output.observation.samples, [12, 12, 12]); + assert.equal(pythonResult.ok, true); + assert.equal(pythonResult.results[0].output.observation.probeId, nodeResult.results[0].output.observation.probeId); + assert.deepEqual(pythonResult.results[0].output.observation.samples, nodeResult.results[0].output.observation.samples); + assert.deepEqual(pythonResult.results[0].output.observation.stats, nodeResult.results[0].output.observation.stats); + assert.equal(pythonResult.results[0].output.observation.rawArtifactRef, nodeResult.results[0].output.observation.rawArtifactRef); const validationPlan = caseValidationPlanFromDefinition({ validation: { mode: "io-probe", ioProbe: { probeId: "main41.ai0.current", quantity: "current", unit: "mA" } } }); const normalized = normalizeValidationStepResult({ step: validationPlan.steps[2], document, exitCode: 0, payload: { body: nodeResult } }); diff --git a/tools/hwlab-node.py b/tools/hwlab-node.py index 173d8d47..0b32b99f 100644 --- a/tools/hwlab-node.py +++ b/tools/hwlab-node.py @@ -615,7 +615,110 @@ class NodeOpsExecutor: argv = args.get("argv") if isinstance(args.get("argv"), list) else [] if not command: return {"ok": False, "blockerCode": "invalid_command", "stderr": "command is required", "exitCode": None} - return self._spawn_output([str(command), *[str(item) for item in argv]], self._workspace_root(args), safe_int(args.get("timeoutMs"), default_timeout_ms)) + command_binding = args.get("commandBinding") if isinstance(args.get("commandBinding"), dict) else {} + if command_binding.get("nativeHelper") == "io-probe-observation": + return self._io_probe_observation(args, command_binding, default_timeout_ms) + output = self._spawn_output([str(command), *[str(item) for item in argv]], self._workspace_root(args), safe_int(args.get("timeoutMs"), default_timeout_ms)) + observation = self._command_observation(command_binding, output) + return {**output, **({"observation": observation} if observation else {})} + + def _command_observation(self, command_binding: dict, output: dict) -> dict: + if command_binding.get("kind") not in ("io-probe", "board-comm"): + return {} + return self._parse_json(output.get("stdout")) + + def _io_probe_observation(self, args: dict, command_binding: dict, default_timeout_ms: int) -> dict: + argv = args.get("argv") if isinstance(args.get("argv"), list) else [] + input_data = self._parse_json(argv[1] if len(argv) > 1 else None) + command_base = input_data.get("commandBase") if isinstance(input_data.get("commandBase"), list) else [] + if not command_base: + return {"ok": False, "blockerCode": "invalid_io_probe_binding", "stderr": "io-probe commandBase is required", "exitCode": None} + count = max(1, min(safe_int(input_data.get("count"), 1), 50)) + interval_ms = self._nonnegative_int(input_data.get("intervalMs")) + settle_ms = self._nonnegative_int(input_data.get("settleMs")) + timeout_ms = safe_int(args.get("timeoutMs"), default_timeout_ms) + samples = [] + response = None + exit_code = 0 + commands = [] + if settle_ms: + time.sleep(settle_ms / 1000) + for index in range(count): + command = self._io_probe_command(input_data, command_base) + output = self._spawn_output(command, self._workspace_root(args), timeout_ms) + commands.append(output) + response = self._parse_json(output.get("stdout")) + value = self._io_probe_value(response, str(input_data.get("valuePath") or "")) + if value is not None: + samples.append(value) + if not output.get("ok"): + exit_code = output.get("exitCode") or 1 + if index < count - 1 and interval_ms: + time.sleep(interval_ms / 1000) + stats = {} + if samples: + stats = {"count": len(samples), "mean": sum(samples) / len(samples), "min": min(samples), "max": max(samples)} + observation = { + "observationId": f"obs_{input_data.get('probeId')}_{int(time.time() * 1000)}", + "probeId": input_data.get("probeId") or command_binding.get("probeId"), + "quantity": input_data.get("quantity") or command_binding.get("quantity"), + "unit": input_data.get("unit") or command_binding.get("unit"), + "samples": samples, + "stats": stats, + "response": response, + "rawArtifactRef": response.get("log_path") or response.get("logPath") or response.get("artifactRef") if isinstance(response, dict) else None, + } + final_exit_code = exit_code if exit_code or samples else 1 + return { + "ok": final_exit_code == 0, + "command": [str(item) for item in command_base], + "cwd": str(self._workspace_root(args)), + "exitCode": final_exit_code, + "stdout": json.dumps(observation, ensure_ascii=False) + "\n", + "stderr": "".join(str(item.get("stderr") or "") for item in commands), + "observation": observation, + } + + def _io_probe_command(self, input_data: dict, command_base: list) -> list[str]: + return [ + *[str(item) for item in command_base], + "jrpctcp", + "--host", + str(input_data.get("host") or ""), + "--port", + str(input_data.get("port") or 8000), + *(["--transport", str(input_data["transport"])] if input_data.get("transport") else []), + *(["--target-node-id", str(input_data["targetNodeId"])] if input_data.get("targetNodeId") is not None else []), + str(input_data.get("method") or "get"), + str(input_data.get("path") or ""), + *[str(item) for item in input_data.get("params", [])], + ] + + def _io_probe_value(self, response: dict, value_path: str) -> float | None: + candidates = [] + nested_response = response.get("response") if isinstance(response, dict) else None + if isinstance(nested_response, dict) and isinstance(nested_response.get("response"), dict): + candidates.append(nested_response["response"]) + if isinstance(nested_response, dict): + candidates.append(nested_response) + if isinstance(response, dict): + candidates.append(response) + keys = [item for item in value_path.removeprefix("$").removeprefix(".").split(".") if item] + for candidate in candidates: + value = candidate + for key in keys: + value = value.get(key) if isinstance(value, dict) else None + try: + return float(value) + except (TypeError, ValueError): + continue + return None + + def _nonnegative_int(self, value: object) -> int: + try: + return max(0, int(value or 0)) + except (TypeError, ValueError): + return 0 def _debug_command(self, op: str, args: dict) -> dict: runs = args.get("commandRuns") if isinstance(args.get("commandRuns"), list) else None