Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
167279e
Prototype LangChain execution-scoped context propagation (util and Op…
lmolkova Sep 30, 2026
e7e7a3c
fix(openai): suspend Responses API stream invocations between reads
sangkyoonnam Sep 30, 2026
29a4afd
docs: describe the stream execution context hook and add changelog fr…
sangkyoonnam Sep 30, 2026
b493216
docs: name the changelog fragments after the pull request
sangkyoonnam Sep 30, 2026
6fcd7b9
Prototype LangChain execution-scoped context propagation (LangChain p…
lmolkova Sep 30, 2026
661226f
fix(langchain): activate context around Runnable._call_with_config
sangkyoonnam Sep 30, 2026
c8992e6
fix(langchain): skip missing execution boundaries when instrumenting
sangkyoonnam Sep 30, 2026
68b988e
test(langchain): cover Responses API stream consumers
sangkyoonnam Sep 30, 2026
4665d51
fix(langchain): attach the parent context where composites enter a step
sangkyoonnam Sep 30, 2026
f34d609
fix(langchain): read the stream config on the first advancement
sangkyoonnam Sep 30, 2026
cb5c1a4
fix(langchain): warn when an installed library lacks an execution bou…
sangkyoonnam Sep 30, 2026
574c447
docs(langchain): list where context is propagated and the remaining gaps
sangkyoonnam Sep 30, 2026
be1f8c5
fix(langchain): find the LangGraph node boundary in pre-0.6 releases
sangkyoonnam Sep 30, 2026
f5e5b7a
chore(langchain): require langchain 0.3.22 and test its langchain-cor…
sangkyoonnam Sep 30, 2026
ecb53ca
fix(langchain): attach the parent context for branch runs and single-…
sangkyoonnam Sep 30, 2026
299e215
fix(langchain): detect the installed library by its distribution befo…
sangkyoonnam Sep 30, 2026
bb030a5
chore(langchain): test LangGraph 0.3.18 as the oldest supported release
sangkyoonnam Sep 30, 2026
54d6b18
fix(langchain): skip the run scope attach when its context is already…
sangkyoonnam Sep 30, 2026
f055087
docs(langchain): explain the patch targets, the run id choice and the…
sangkyoonnam Sep 30, 2026
e1d9690
test(langchain): pin the wrap order and restoration of the graph stre…
sangkyoonnam Sep 30, 2026
1aa6a66
fix(langchain): end the chat span when ainvoke is cancelled or interr…
sangkyoonnam Sep 30, 2026
02612b0
docs(langchain): name the changelog fragments after the pull request
sangkyoonnam Sep 30, 2026
35867a4
fix(util): run the stream manager exit inside the stream execution co…
sangkyoonnam Sep 30, 2026
b2538dc
fix(openai): finalize a close through the Responses API HTTP response…
sangkyoonnam Sep 30, 2026
4253242
docs(util): note the execution context on close finalizers and manage…
sangkyoonnam Sep 30, 2026
4ef42d0
Merge branch 'main' into fix/677-langchain-execution-context-split
sangkyoonnam Oct 6, 2026
8b3015b
fix(langchain): mint run ids with LangChain's uuid7 when available
sangkyoonnam Oct 6, 2026
db80ddf
test(langchain): skip async composite cases on Python 3.10 with langc…
sangkyoonnam Oct 6, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion .github/instructions/instrumentation.instructions.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,9 @@ prefer opt-in or additive. Breaking changes need explicit justification in the P
capture path — never as unconditional span/log attributes.
- Adding attributes to invocations produced by the util is fine.
- Streaming responses must be instrumented by subclassing the util's `SyncStreamWrapper` /
`AsyncStreamWrapper` (`opentelemetry.util.genai.stream`). Flag hand-rolled stream wrappers.
`AsyncStreamWrapper` (`opentelemetry.util.genai.stream`). Flag hand-rolled stream wrappers, and
invocation-backed wrappers that do not `invocation.suspend()` before the stream is returned and
return `invocation.activate()` from `_execution_context()`.
- Instrumentation should not change what a call returns or when its work happens. Flag: work the SDK
didn't do (materializing a result early to build telemetry — stay lazy); a changed return type
(`isinstance`/`__class__` should still resolve to the original; `wrapt.ObjectProxy` is the usual
Expand Down
14 changes: 12 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -229,14 +229,20 @@ as the reference:

A streamed response only finishes once the caller has drained the stream, so the invocation must
stay open until then. Do **not** call `invocation.stop()` when the SDK returns the stream — the
span would close before any chunks arrive.
span would close before any chunks arrive. The invocation's span must also not stay current in
the caller's context while the stream is unconsumed: call `invocation.suspend()` before returning
the stream, and re-activate the invocation only while a chunk is being read.

Instrument streams by subclassing `SyncStreamWrapper` / `AsyncStreamWrapper` from
`opentelemetry.util.genai.stream` (the public, supported helpers). The base class proxies the
underlying SDK stream, drives iteration, and finalizes telemetry exactly once on success, error,
or `close()`. Subclasses pass the SDK stream to `super().__init__(stream)` and implement three
or `close()`. Subclasses pass the SDK stream to `super().__init__(stream)` and implement four
hooks:

- `_execution_context()` — return a fresh context manager for each stream read and cleanup
operation; for invocation-backed streams return `invocation.activate()`. Context is restored
before the chunk is returned to the consumer, so the token is created and released in the
same frame.
- `_process_chunk(chunk)` — accumulate per-chunk state (e.g. response model, finish reasons,
token usage, streamed content) onto the invocation.
- `_on_stream_end()` — finalize on success; set the accumulated response attributes and call
Expand All @@ -248,8 +254,12 @@ class MyStreamWrapper(SyncStreamWrapper[Chunk]):
def __init__(self, stream, invocation, capture_content):
super().__init__(stream)
self._self_invocation = invocation
invocation.suspend()
...

def _execution_context(self):
return self._self_invocation.activate()

def _process_chunk(self, chunk): ... # accumulate state
def _on_stream_end(self):
self._self_invocation.stop()
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ All instrumentations use [opentelemetry-util-genai](./util/opentelemetry-util-ge
| [opentelemetry-instrumentation-genai-anthropic](./instrumentation/opentelemetry-instrumentation-genai-anthropic) | anthropic >= 0.51.0, < 2 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-anthropic/) |
| [opentelemetry-instrumentation-genai-bedrock](./instrumentation/opentelemetry-instrumentation-genai-bedrock) | boto3 >= 1.40.46, < 2 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-bedrock/) |
| [opentelemetry-instrumentation-genai-dspy](./instrumentation/opentelemetry-instrumentation-genai-dspy) | dspy >= 3.3.0, < 4 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-dspy/) |
| [opentelemetry-instrumentation-genai-langchain](./instrumentation/opentelemetry-instrumentation-genai-langchain) | langchain >= 0.3.21, < 2 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-langchain/) |
| [opentelemetry-instrumentation-genai-langchain](./instrumentation/opentelemetry-instrumentation-genai-langchain) | langchain >= 0.3.22, < 2 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-langchain/) |
| [opentelemetry-instrumentation-genai-llama-index](./instrumentation/opentelemetry-instrumentation-genai-llama-index) | llama-index-core >= 0.14.19, < 1, llama-index-instrumentation >= 0.4.3, < 1, llama-index-workflows >= 2.17.1, != 2.24.0, < 3 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-llama-index/) |
| [opentelemetry-instrumentation-genai-openai](./instrumentation/opentelemetry-instrumentation-genai-openai) | openai >= 1.26.0, < 4 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-openai/) |
| [opentelemetry-instrumentation-genai-openai-agents](./instrumentation/opentelemetry-instrumentation-genai-openai-agents) | openai-agents >= 0.3.3, < 1 | [1.2b0](https://pypi.org/project/opentelemetry-instrumentation-genai-openai-agents/) |
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Require `langchain >= 0.3.22`: the execution boundaries rely on `set_config_context`, which langchain-core added in 0.3.46 (langchain 0.3.22 is the first release whose langchain-core floor includes it). LangGraph, when installed, needs 0.3.18 or newer for node boundaries.
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Activate the span context around LangChain and LangGraph execution boundaries (chat models, tools, retrievers, runnables, streams, graph nodes) instead of attaching it from callbacks, so nested SDK and HTTP spans correlate across asyncio tasks and threads without `Failed to detach context` errors. A chat model call cancelled or interrupted during `ainvoke` now ends its span with that error instead of leaking it.
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@ Installation

pip install opentelemetry-instrumentation-genai-langchain

LangGraph is optional. Graph node boundaries are instrumented from LangGraph
0.3.18 on; an older LangGraph still runs, with a warning that context is not
propagated across its nodes.

See the `examples <examples>`_ directory for runnable ``workflow``, ``agent``,
``tools``, and ``zero-code`` scenarios.

Expand Down Expand Up @@ -141,8 +145,14 @@ programmatically, which takes precedence over the environment variable::
Known Limitations
-----------------

Context propagation to nested calls (such as auto-instrumented HTTP clients
or database queries within tools) is not supported when using LangChain async API.
Context is propagated to nested calls (such as auto-instrumented HTTP clients or
database queries) where LangChain hands control to user code: chat models, tools,
retrievers, runnables built on ``_call_with_config``, streams, the steps of sequence,
parallel, branch and fallback runnables, single-input ``batch``/``abatch`` calls, and
LangGraph nodes. A ``Runnable`` whose ``invoke`` starts no run of its own is not
correlated when ``batch``/``abatch`` runs it for several inputs at once: the default
implementation dispatches each input from a thread pool or task with no frame of its
own to carry that input's parent.

References
----------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,12 @@ classifiers = [
]
dependencies = [
"opentelemetry-instrumentation >= 0.64b0, <1",
"opentelemetry-util-genai >= 1.2b0, <2",
"opentelemetry-util-genai >= 1.3b0.dev, <2",
]

[project.optional-dependencies]
instruments = [
"langchain >= 0.3.21, < 2",
"langchain >= 0.3.22, < 2",
]

[project.entry-points.opentelemetry_instrumentor]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,11 @@
from typing import Any

from langchain_core.callbacks import BaseCallbackManager
from langchain_core.callbacks.manager import AsyncCallbackManager
from wrapt import wrap_function_wrapper

from opentelemetry.instrumentation.genai.langchain._execution_context import (
_ExecutionContext,
)
from opentelemetry.instrumentation.genai.langchain.agent_context import (
wrap_astream,
wrap_stream,
Expand All @@ -57,6 +59,8 @@ class LangChainInstrumentor(BaseInstrumentor):
to capture LLM telemetry.
"""

_execution_context: _ExecutionContext | None = None

def __init__(
self,
):
Expand All @@ -83,22 +87,22 @@ def _instrument(self, **kwargs: Any):
instrumentation_scope_version=__version__,
)
invocation_manager = _InvocationManager()
sync_handler = OpenTelemetryLangChainCallbackHandler(
telemetry_handler=telemetry_handler,
_attach_to_context=True,
invocation_manager=invocation_manager,
)
async_handler = OpenTelemetryLangChainCallbackHandler(
handler = OpenTelemetryLangChainCallbackHandler(
telemetry_handler=telemetry_handler,
_attach_to_context=False,
invocation_manager=invocation_manager,
)

wrap_function_wrapper(
"langchain_core.callbacks",
"BaseCallbackManager.__init__",
_BaseCallbackManagerInitWrapper(sync_handler, async_handler),
_BaseCallbackManagerInitWrapper(handler),
)
self._execution_context = _ExecutionContext(
invocation_manager,
handler.on_tool_error,
handler.on_retriever_error,
)
self._execution_context.instrument()
self._instrument_agent_entry_points()

@staticmethod
Expand All @@ -120,13 +124,18 @@ def _uninstrument(self, **kwargs: Any):
Cleanup instrumentation (unwrap).
"""
unwrap("langchain_core.callbacks.base.BaseCallbackManager", "__init__")
# The agent entry points wrap last, over the execution boundary, so
# they come off first.
try:
import langgraph.pregel

for method in ("stream", "astream"):
unwrap(langgraph.pregel.Pregel, method)
except (ImportError, AttributeError):
pass
if self._execution_context is not None:
self._execution_context.uninstrument()
self._execution_context = None


class _BaseCallbackManagerInitWrapper:
Expand All @@ -136,11 +145,9 @@ class _BaseCallbackManagerInitWrapper:

def __init__(
self,
sync_handler: OpenTelemetryLangChainCallbackHandler,
async_handler: OpenTelemetryLangChainCallbackHandler,
handler: OpenTelemetryLangChainCallbackHandler,
):
self._sync_handler = sync_handler
self._async_handler = async_handler
self._handler = handler

def __call__(
self,
Expand All @@ -150,29 +157,8 @@ def __call__(
kwargs: dict[str, Any],
):
wrapped(*args, **kwargs)
target_handler = (
self._async_handler
if isinstance(instance, AsyncCallbackManager)
else self._sync_handler
)
other_handler = (
self._sync_handler
if isinstance(instance, AsyncCallbackManager)
else self._async_handler
)
if other_handler in instance.handlers:
instance.handlers = [
target_handler if h is other_handler else h
for h in instance.handlers
]
if other_handler in instance.inheritable_handlers:
instance.inheritable_handlers = [
target_handler if h is other_handler else h
for h in instance.inheritable_handlers
]
if target_handler not in instance.inheritable_handlers:
for handler in instance.inheritable_handlers:
if isinstance(handler, OpenTelemetryLangChainCallbackHandler):
break
else:
instance.add_handler(target_handler, inherit=True)
if not any(
isinstance(handler, OpenTelemetryLangChainCallbackHandler)
for handler in instance.inheritable_handlers
):
instance.add_handler(self._handler, inherit=True)
Loading