From 58e0773ec6c02ea1aca499e934b269dbc55f812c Mon Sep 17 00:00:00 2001 From: Nick Hollon Date: Fri, 8 May 2026 11:38:27 -0400 Subject: [PATCH] rename trigger_call_id to parent_task_id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pregel doesn't have a "call" concept — its unit of work is a task. The namespace segment after `:` is documented in `_internal/_constants.py` as the `task_id` ("for checkpoint_ns, for each level, separates the namespace from the task_id"), and "parent task" already shows up as existing pregel vocabulary (`_algo.py`, `_runner.py`). `parent_task_id` matches that vocabulary: from the child subgraph's perspective, the id we carry on every lifecycle event identifies the task in the parent scope whose execution dispatched this subgraph. Renames the wire field, the `SubgraphRunStream.parent_task_id` attribute, all internal parameter / variable names, and the test identifiers. --- libs/langgraph/langgraph/stream/run_stream.py | 10 +-- .../langgraph/stream/transformers.py | 78 +++++++++---------- .../test_stream_lifecycle_transformer.py | 28 +++---- .../tests/test_stream_subgraph_transformer.py | 2 +- 4 files changed, 59 insertions(+), 59 deletions(-) diff --git a/libs/langgraph/langgraph/stream/run_stream.py b/libs/langgraph/langgraph/stream/run_stream.py index d0126a22f..33b3bd48d 100644 --- a/libs/langgraph/langgraph/stream/run_stream.py +++ b/libs/langgraph/langgraph/stream/run_stream.py @@ -523,7 +523,7 @@ class _SubgraphRunStreamMixin: path: tuple[str, ...] graph_name: str | None - trigger_call_id: str | None + parent_task_id: str | None status: LifecycleEvent error: str | None _seen_terminal: bool @@ -538,7 +538,7 @@ class SubgraphRunStream(GraphRunStream, _SubgraphRunStreamMixin): *, path: tuple[str, ...], graph_name: str | None = None, - trigger_call_id: str | None = None, + parent_task_id: str | None = None, ) -> None: # Capture the parent-inherited pump before super().__init__ # touches anything; we delegate to it from `_pump_next`. @@ -550,7 +550,7 @@ class SubgraphRunStream(GraphRunStream, _SubgraphRunStreamMixin): ) self.path = path self.graph_name = graph_name - self.trigger_call_id = trigger_call_id + self.parent_task_id = parent_task_id self.status = "started" self.error = None self._seen_terminal = False @@ -581,7 +581,7 @@ class AsyncSubgraphRunStream(AsyncGraphRunStream, _SubgraphRunStreamMixin): *, path: tuple[str, ...], graph_name: str | None = None, - trigger_call_id: str | None = None, + parent_task_id: str | None = None, ) -> None: self._parent_apump_fn: Callable[[], Awaitable[bool]] | None = mux._apump_fn super().__init__( @@ -591,7 +591,7 @@ class AsyncSubgraphRunStream(AsyncGraphRunStream, _SubgraphRunStreamMixin): ) self.path = path self.graph_name = graph_name - self.trigger_call_id = trigger_call_id + self.parent_task_id = parent_task_id self.status = "started" self.error = None self._seen_terminal = False diff --git a/libs/langgraph/langgraph/stream/transformers.py b/libs/langgraph/langgraph/stream/transformers.py index 79f837372..605581c54 100644 --- a/libs/langgraph/langgraph/stream/transformers.py +++ b/libs/langgraph/langgraph/stream/transformers.py @@ -333,7 +333,7 @@ LifecycleEvent = Literal["started", "completed", "failed", "interrupted", "drain Each value: - `started` — first `tasks` event observed at the tracked namespace. Carries - `graph_name` / `trigger_call_id` / optional `cause` describing what spawned + `graph_name` / `parent_task_id` / optional `cause` describing what spawned the subgraph. - `completed` — the dispatching task's `TaskResultPayload` arrived with neither error nor interrupts. Also emitted by `finalize` for any tracked namespace @@ -351,7 +351,7 @@ Each value: def _parse_ns_segment(segment: str) -> tuple[str, str | None]: - """Split a namespace segment into `(graph_name, trigger_call_id)`. + """Split a namespace segment into `(graph_name, parent_task_id)`. Segments are formatted `node_name:task_id` by `prepare_next_tasks`. Returns `(segment, None)` if no `:` is present. @@ -416,7 +416,7 @@ class LifecyclePayload(TypedDict, total=False): split (it does for normally-dispatched subgraphs). Absent if the segment has no `:` separator. """ - trigger_call_id: str + parent_task_id: str """Pregel task id of the dispatching task — the task whose execution spawned this subgraph. @@ -465,10 +465,10 @@ class _TasksLifecycleBase(StreamTransformer): - `_should_track(ns)` — scope filter (e.g. multi-depth vs direct-children-only). - - `_on_started(ns, graph_name, trigger_call_id, tool_call_id)` — + - `_on_started(ns, graph_name, parent_task_id, tool_call_id)` — first sighting action (push payload / build handle / etc.). Called once per discovered namespace. - - `_on_terminal(ns, status, error, trigger_call_id)` — terminal + - `_on_terminal(ns, status, error, parent_task_id)` — terminal action (push terminal payload / mark handle status). Called once per tracked namespace at result time, or via `finalize` / `fail` sweeps if no parent result arrived. @@ -504,7 +504,7 @@ class _TasksLifecycleBase(StreamTransformer): self, ns: tuple[str, ...], graph_name: str | None, - trigger_call_id: str | None, + parent_task_id: str | None, tool_call_id: str | None = None, ) -> None: """Fired once per discovered namespace (first observed task event). @@ -512,7 +512,7 @@ class _TasksLifecycleBase(StreamTransformer): `tool_call_id` is the model-side id of the originating tool call (from the per-call dispatched task's `input`). `None` for structurally-triggered subgraphs or per-call envelopes that omitted - an `id`. Consumers join on `trigger_call_id` (the pregel task id) + an `id`. Consumers join on `parent_task_id` (the pregel task id) for identity; `tool_call_id` is purely an anchor back to the AI message that dispatched the subgraph. """ @@ -523,12 +523,12 @@ class _TasksLifecycleBase(StreamTransformer): ns: tuple[str, ...], status: LifecycleEvent, error: str | None, - trigger_call_id: str, + parent_task_id: str, ) -> None: """Fired once per tracked namespace when its dispatching task's result arrives, or via finalize/fail safety-net sweeps. - `trigger_call_id` is the same id paired with the namespace at + `parent_task_id` is the same id paired with the namespace at `_on_started` time, so subscribers can correlate the terminal event back to its `started`. """ @@ -553,26 +553,26 @@ class _TasksLifecycleBase(StreamTransformer): def _handle_task_start(self, ns: tuple[str, ...], data: dict[str, Any]) -> None: # Mine input shape on every tasks event (not just tracked ones) # so we capture dispatching tasks that themselves live outside the - # tracked region but whose `id` will appear as `trigger_call_id` + # tracked region but whose `id` will appear as `parent_task_id` # for a child subgraph. self._record_dispatching_tool_call_id(data) if not self._should_track(ns) or ns in self._seen: return self._seen.add(ns) - graph_name, trigger_call_id = _parse_ns_segment(ns[-1]) + graph_name, parent_task_id = _parse_ns_segment(ns[-1]) tool_call_id = ( - self._dispatching_tool_call_id.pop(trigger_call_id, None) - if trigger_call_id is not None + self._dispatching_tool_call_id.pop(parent_task_id, None) + if parent_task_id is not None else None ) self._on_started( ns, graph_name or None, - trigger_call_id, + parent_task_id, tool_call_id, ) - if trigger_call_id is not None: - self._open[ns] = trigger_call_id + if parent_task_id is not None: + self._open[ns] = parent_task_id def _record_dispatching_tool_call_id(self, data: dict[str, Any]) -> None: """Remember `task_id -> tool_call_id` if the task input matches @@ -597,8 +597,8 @@ class _TasksLifecycleBase(StreamTransformer): ) -> list[tuple[tuple[str, ...], LifecycleEvent, str | None, str]]: """Return and remove tracked children closed by this task result. - Each tuple is `(child_ns, status, error, trigger_call_id)`. - `trigger_call_id` is the dispatching task's id — the same id + Each tuple is `(child_ns, status, error, parent_task_id)`. + `parent_task_id` is the dispatching task's id — the same id we'd already paired with the namespace at `_on_started`. """ result_id = data.get("id") @@ -614,22 +614,22 @@ class _TasksLifecycleBase(StreamTransformer): return transitions def _handle_task_result(self, ns: tuple[str, ...], data: dict[str, Any]) -> None: - for child_ns, status, error, trigger_call_id in self._pop_terminal_transitions( + for child_ns, status, error, parent_task_id in self._pop_terminal_transitions( ns, data ): - self._on_terminal(child_ns, status, error, trigger_call_id) + self._on_terminal(child_ns, status, error, parent_task_id) def finalize(self) -> None: """Emit `completed` for any tracked namespace still open at run end.""" - for ns, trigger_call_id in list(self._open.items()): - self._on_terminal(ns, "completed", None, trigger_call_id) + for ns, parent_task_id in list(self._open.items()): + self._on_terminal(ns, "completed", None, parent_task_id) self._open.clear() def fail(self, err: BaseException) -> None: """Emit terminal status for any tracked namespace still open.""" status, error_str = _status_from_exception(err) - for ns, trigger_call_id in list(self._open.items()): - self._on_terminal(ns, status, error_str, trigger_call_id) + for ns, parent_task_id in list(self._open.items()): + self._on_terminal(ns, status, error_str, parent_task_id) self._open.clear() @@ -693,10 +693,10 @@ class LifecycleTransformer(_TasksLifecycleBase): self, ns: tuple[str, ...], graph_name: str | None, - trigger_call_id: str | None, + parent_task_id: str | None, tool_call_id: str | None = None, ) -> None: - if trigger_call_id is None: + if parent_task_id is None: # Without a task id we can't correlate a dispatching-task-result # event back to this namespace — skip the started payload # and rely on finalize/fail to close. @@ -704,7 +704,7 @@ class LifecycleTransformer(_TasksLifecycleBase): payload: LifecyclePayload = { "event": "started", "namespace": list(ns), - "trigger_call_id": trigger_call_id, + "parent_task_id": parent_task_id, } if graph_name: payload["graph_name"] = graph_name @@ -717,12 +717,12 @@ class LifecycleTransformer(_TasksLifecycleBase): ns: tuple[str, ...], status: LifecycleEvent, error: str | None, - trigger_call_id: str, + parent_task_id: str, ) -> None: payload: LifecyclePayload = { "event": status, "namespace": list(ns), - "trigger_call_id": trigger_call_id, + "parent_task_id": parent_task_id, } if error is not None: payload["error"] = error @@ -776,7 +776,7 @@ class SubgraphTransformer(_TasksLifecycleBase): self, ns: tuple[str, ...], graph_name: str | None, - trigger_call_id: str | None, + parent_task_id: str | None, tool_call_id: str | None = None, # noqa: ARG002 ) -> None: if self._mux is None: @@ -790,7 +790,7 @@ class SubgraphTransformer(_TasksLifecycleBase): mux=child_mux, path=ns, graph_name=graph_name, - trigger_call_id=trigger_call_id, + parent_task_id=parent_task_id, ) self._handles[ns] = handle self._log.push(handle) @@ -800,7 +800,7 @@ class SubgraphTransformer(_TasksLifecycleBase): ns: tuple[str, ...], status: LifecycleEvent, error: str | None, - trigger_call_id: str, # noqa: ARG002 + parent_task_id: str, # noqa: ARG002 ) -> None: handle = self._handles.get(ns) if handle is None or not self._mark_terminal(handle, status, error): @@ -812,7 +812,7 @@ class SubgraphTransformer(_TasksLifecycleBase): ns: tuple[str, ...], status: LifecycleEvent, error: str | None, - trigger_call_id: str, # noqa: ARG002 + parent_task_id: str, # noqa: ARG002 ) -> None: handle = self._handles.get(ns) if handle is None or not self._mark_terminal(handle, status, error): @@ -893,9 +893,9 @@ class SubgraphTransformer(_TasksLifecycleBase): child_ns, status, error, - trigger_call_id, + parent_task_id, ) in self._pop_terminal_transitions(ns, data): - await self._aon_terminal(child_ns, status, error, trigger_call_id) + await self._aon_terminal(child_ns, status, error, parent_task_id) else: self._handle_task_start(ns, data) keep = False @@ -909,9 +909,9 @@ class SubgraphTransformer(_TasksLifecycleBase): def _complete_open_handles(self) -> BaseException | None: first_error: BaseException | None = None - for ns, trigger_call_id in list(self._open.items()): + for ns, parent_task_id in list(self._open.items()): try: - self._on_terminal(ns, "completed", None, trigger_call_id) + self._on_terminal(ns, "completed", None, parent_task_id) except BaseException as e: if first_error is None: first_error = e @@ -927,9 +927,9 @@ class SubgraphTransformer(_TasksLifecycleBase): async def _acomplete_open_handles(self) -> BaseException | None: first_error: BaseException | None = None - for ns, trigger_call_id in list(self._open.items()): + for ns, parent_task_id in list(self._open.items()): try: - await self._aon_terminal(ns, "completed", None, trigger_call_id) + await self._aon_terminal(ns, "completed", None, parent_task_id) except BaseException as e: if first_error is None: first_error = e diff --git a/libs/langgraph/tests/test_stream_lifecycle_transformer.py b/libs/langgraph/tests/test_stream_lifecycle_transformer.py index ef72a187a..fe38c96c9 100644 --- a/libs/langgraph/tests/test_stream_lifecycle_transformer.py +++ b/libs/langgraph/tests/test_stream_lifecycle_transformer.py @@ -131,7 +131,7 @@ def test_started_emitted_on_first_direct_child_task() -> None: assert payload["event"] == "started" assert payload["namespace"] == ["agent:abc123"] assert payload["graph_name"] == "agent" - assert payload["trigger_call_id"] == "abc123" + assert payload["parent_task_id"] == "abc123" def test_started_carries_cause_for_dict_envelope_input() -> None: @@ -142,7 +142,7 @@ def test_started_carries_cause_for_dict_envelope_input() -> None: When that task triggers a subgraph (the child's namespace ends in `name:`), `lifecycle.started.cause` carries `{"type": "tool_call", "tool_call_id": ...}`. Identity correlation - still uses `trigger_call_id`; `tool_call_id` is exposed so UI + still uses `parent_task_id`; `tool_call_id` is exposed so UI consumers can anchor the lifecycle event back to the originating AI message tool call. Args are deliberately NOT mined — they live on the AIMessage and have a single source of truth there. @@ -166,7 +166,7 @@ def test_started_carries_cause_for_dict_envelope_input() -> None: [payload] = _drain_lifecycle(mux) assert payload["event"] == "started" - assert payload["trigger_call_id"] == "abc123" + assert payload["parent_task_id"] == "abc123" assert payload["cause"] == { "type": "tool_call", "tool_call_id": "call_xyz", @@ -200,7 +200,7 @@ def test_started_carries_cause_for_list_shape_per_call_input() -> None: [payload] = _drain_lifecycle(mux) assert payload["event"] == "started" - assert payload["trigger_call_id"] == "abc123" + assert payload["parent_task_id"] == "abc123" assert payload["cause"] == {"type": "tool_call", "tool_call_id": "tc-1"} @@ -284,7 +284,7 @@ def test_parallel_dispatches_attributed_to_correct_parent() -> None: out to their own child subgraph; each child's `cause.tool_call_id` must reflect its own dispatching envelope, not the other. - Defends the `trigger_call_id` (pregel task id) join: that id is + Defends the `parent_task_id` (pregel task id) join: that id is parsed from the child namespace segment and is unique per Send, so it disambiguates parallel dispatches 1:1. Both children share the same `subagent_type` (in args, not on cause) — only the @@ -411,10 +411,10 @@ def test_completed_on_parent_task_result() -> None: payloads = _drain_lifecycle(mux) assert [p["event"] for p in payloads] == ["started", "completed"] - # `trigger_call_id` is required on every event for the same subgraph + # `parent_task_id` is required on every event for the same subgraph # so consumers can correlate `started` ↔ terminal without joining # via `namespace`. - assert all(p["trigger_call_id"] == "abc" for p in payloads) + assert all(p["parent_task_id"] == "abc" for p in payloads) def test_failed_on_parent_task_result_with_error() -> None: @@ -490,10 +490,10 @@ def test_fail_emits_failed_for_other_exceptions() -> None: assert payloads[1]["error"] == "boom" -def test_trigger_call_id_present_on_every_terminal_path() -> None: +def test_parent_task_id_present_on_every_terminal_path() -> None: """Every exit path that emits a terminal event (parent-result with error / interrupts, finalize sweep, fail sweep) must carry - `trigger_call_id` so consumers can correlate the terminal event + `parent_task_id` so consumers can correlate the terminal event back to its `started` without falling back to namespace joins.""" # Path 1: parent-result with error. mux = _build_lifecycle_mux() @@ -501,7 +501,7 @@ def test_trigger_call_id_present_on_every_terminal_path() -> None: mux.push(_tasks_result([], task_id="abc", name="agent", error="boom")) [_, terminal] = _drain_lifecycle(mux) assert terminal["event"] == "failed" - assert terminal["trigger_call_id"] == "abc" + assert terminal["parent_task_id"] == "abc" # Path 2: parent-result with interrupts. mux2 = _build_lifecycle_mux() @@ -511,7 +511,7 @@ def test_trigger_call_id_present_on_every_terminal_path() -> None: ) [_, terminal2] = _drain_lifecycle(mux2) assert terminal2["event"] == "interrupted" - assert terminal2["trigger_call_id"] == "def" + assert terminal2["parent_task_id"] == "def" # Path 3: finalize sweep (no parent result arrived). mux3 = _build_lifecycle_mux() @@ -519,7 +519,7 @@ def test_trigger_call_id_present_on_every_terminal_path() -> None: mux3.close() [_, terminal3] = _drain_lifecycle(mux3) assert terminal3["event"] == "completed" - assert terminal3["trigger_call_id"] == "ghi" + assert terminal3["parent_task_id"] == "ghi" # Path 4: fail sweep with GraphInterrupt. mux4 = _build_lifecycle_mux() @@ -527,7 +527,7 @@ def test_trigger_call_id_present_on_every_terminal_path() -> None: mux4.fail(GraphInterrupt()) [_, terminal4] = _drain_lifecycle(mux4) assert terminal4["event"] == "interrupted" - assert terminal4["trigger_call_id"] == "jkl" + assert terminal4["parent_task_id"] == "jkl" # Path 5: fail sweep with generic exception. mux5 = _build_lifecycle_mux() @@ -535,7 +535,7 @@ def test_trigger_call_id_present_on_every_terminal_path() -> None: mux5.fail(RuntimeError("kaboom")) [_, terminal5] = _drain_lifecycle(mux5) assert terminal5["event"] == "failed" - assert terminal5["trigger_call_id"] == "mno" + assert terminal5["parent_task_id"] == "mno" def test_unrelated_methods_pass_through() -> None: diff --git a/libs/langgraph/tests/test_stream_subgraph_transformer.py b/libs/langgraph/tests/test_stream_subgraph_transformer.py index 236f73257..ee0ad7634 100644 --- a/libs/langgraph/tests/test_stream_subgraph_transformer.py +++ b/libs/langgraph/tests/test_stream_subgraph_transformer.py @@ -195,7 +195,7 @@ def test_handle_created_on_first_direct_child_task() -> None: [handle] = _drain_subgraphs(mux) assert handle.path == ("agent:abc",) assert handle.graph_name == "agent" - assert handle.trigger_call_id == "abc" + assert handle.parent_task_id == "abc" assert handle.status == "started" _child_mux(handle) # mini-mux backed