Run a BaseNode through NodeRunner. Events flow through ic._event_queue via NodeRunner.
(
self,
*,
user_id: str,
session_id: str,
invocation_id: Optional[str] = None,
new_message: Optional[types.Content] = None,
state_delta: Optional[dict[str, Any]] = None,
run_config: Optional[RunConfig] = None,
yield_user_message: bool = False,
node: Optional['BaseNode'] = None,
)
| 435 | ) |
| 436 | |
| 437 | async def _run_node_async( |
| 438 | self, |
| 439 | *, |
| 440 | user_id: str, |
| 441 | session_id: str, |
| 442 | invocation_id: Optional[str] = None, |
| 443 | new_message: Optional[types.Content] = None, |
| 444 | state_delta: Optional[dict[str, Any]] = None, |
| 445 | run_config: Optional[RunConfig] = None, |
| 446 | yield_user_message: bool = False, |
| 447 | node: Optional['BaseNode'] = None, |
| 448 | ) -> AsyncGenerator[Event, None]: |
| 449 | """Run a BaseNode through NodeRunner. |
| 450 | |
| 451 | Events flow through ic._event_queue via NodeRunner. |
| 452 | """ |
| 453 | from .workflow._node_runner import NodeRunner |
| 454 | |
| 455 | with tracer.start_as_current_span('invocation'): |
| 456 | # 1. Setup |
| 457 | session = await self._get_or_create_session( |
| 458 | user_id=user_id, session_id=session_id |
| 459 | ) |
| 460 | |
| 461 | # Validate and resolve resume inputs |
| 462 | resume_inputs = self._extract_resume_inputs(new_message) |
| 463 | self._validate_new_message(new_message, resume_inputs) |
| 464 | |
| 465 | if not invocation_id and new_message: |
| 466 | invocation_id = self._resolve_invocation_id_from_fr( |
| 467 | session, new_message |
| 468 | ) |
| 469 | |
| 470 | ic = self._new_invocation_context( |
| 471 | session, |
| 472 | new_message=new_message, |
| 473 | run_config=run_config or RunConfig(), |
| 474 | invocation_id=invocation_id, |
| 475 | ) |
| 476 | ic._event_queue = asyncio.Queue() |
| 477 | |
| 478 | # 2. Append user message to session and resolve node_input |
| 479 | node_input = None |
| 480 | if resume_inputs or invocation_id: |
| 481 | # Resume: recover the original user content. new_message here is a |
| 482 | # function response (or None), so it can't populate user_content. |
| 483 | node_input = self._find_original_user_content( |
| 484 | ic.session, ic.invocation_id |
| 485 | ) |
| 486 | if node_input: |
| 487 | ic.user_content = node_input |
| 488 | if not node_input: |
| 489 | # Fresh: use user message as node_input |
| 490 | node_input = new_message |
| 491 | |
| 492 | # Run callbacks on user message |
| 493 | if new_message: |
| 494 | modified_user_message = ( |
no test coverage detected