Record a node-start event. Args: node: Node instance (or internal node object) being executed. inputs: Input payload passed to the node. Returns: A `_NodeRunContext` used to later record end/error events, or `None` if the node cannot
(self, node: object, inputs: dict[str, object])
| 328 | return None |
| 329 | |
| 330 | def node_start(self, node: object, inputs: dict[str, object]) -> _NodeRunContext | None: |
| 331 | """Record a node-start event. |
| 332 | |
| 333 | Args: |
| 334 | node: Node instance (or internal node object) being executed. |
| 335 | inputs: Input payload passed to the node. |
| 336 | |
| 337 | Returns: |
| 338 | A `_NodeRunContext` used to later record end/error events, or `None` if the node |
| 339 | cannot be resolved to a stable node id. |
| 340 | """ |
| 341 | node_name = self._resolve_node_id(node) |
| 342 | if not node_name: |
| 343 | return None |
| 344 | |
| 345 | token_before_in = None |
| 346 | token_before_out = None |
| 347 | try: |
| 348 | model = getattr(node, "model", None) |
| 349 | tracker = getattr(model, "token_tracker", None) if model else None |
| 350 | token_before_in = int(getattr(tracker, "total_input_usage", 0)) if tracker else None |
| 351 | token_before_out = int(getattr(tracker, "total_output_usage", 0)) if tracker else None |
| 352 | except Exception: |
| 353 | token_before_in = None |
| 354 | token_before_out = None |
| 355 | |
| 356 | run_id = f"{self._pid}-{uuid.uuid4().hex[:10]}" |
| 357 | ts = _now_ms() |
| 358 | |
| 359 | payload: dict[str, object] = { |
| 360 | "type": "NODE_EVENT", |
| 361 | "event": "start", |
| 362 | "node": node_name, |
| 363 | "ts": ts, |
| 364 | "runId": run_id, |
| 365 | "inputs": self._safe_for_history(inputs if isinstance(inputs, dict) else {"input": inputs}), |
| 366 | } |
| 367 | if self.is_streaming(): |
| 368 | self._enqueue(payload) |
| 369 | else: |
| 370 | self._record_history(payload) |
| 371 | return _NodeRunContext( |
| 372 | node_name=node_name, |
| 373 | run_id=run_id, |
| 374 | started_ms=ts, |
| 375 | token_before_in=token_before_in, |
| 376 | token_before_out=token_before_out, |
| 377 | ) |
| 378 | |
| 379 | def node_end(self, ctx: _NodeRunContext | None, outputs: dict[str, object], node: object | None = None) -> None: |
| 380 | """Record a node-end event and attach basic metrics when available. |
nothing calls this directly
no test coverage detected