Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
33 changes: 3 additions & 30 deletions python/tracing/langchain/00_quickstart.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,48 +9,21 @@

from __future__ import annotations

import os

from dotenv import find_dotenv, load_dotenv
from _shared import init_telemetry, tracing_config
from langchain_core.language_models.fake_chat_models import FakeListChatModel
from langchain_core.messages import HumanMessage, SystemMessage
from respan_instrumentation_langchain import add_respan_callback
from respan_tracing import RespanTelemetry

load_dotenv(find_dotenv(), override=False)


def langchain_instrumentation_quickstart() -> None:
api_key = os.getenv("RESPAN_API_KEY")
telemetry: RespanTelemetry | None = None

if api_key:
telemetry = RespanTelemetry(
app_name="langchain-quickstart",
api_key=api_key,
base_url=os.getenv("RESPAN_BASE_URL", "https://api.respan.ai/api"),
is_auto_instrument=False,
is_batching_enabled=False,
is_enabled=True,
)
else:
print("RESPAN_API_KEY is not set; running locally without exporting spans.")
init_telemetry("langchain-quickstart")

model = FakeListChatModel(responses=["Hello from a traced LangChain run."])
config = {
"run_name": "hello_world",
"tags": ["respan-langchain-example", "quickstart"],
"metadata": {"example": "quickstart"},
}
if telemetry:
config = add_respan_callback(config)

response = model.invoke(
[
SystemMessage(content="Reply in one short sentence."),
HumanMessage(content="Say hello to Respan tracing."),
],
config=config,
config=tracing_config("hello_world"),
)
print(response.content)

Expand Down
9 changes: 5 additions & 4 deletions python/tracing/langchain/06_chat_model_astream.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,20 +2,21 @@

import asyncio

from langchain_core.language_models.fake_chat_models import FakeChatModel

from _shared import init_telemetry, message_text, tracing_config
from langchain_core.language_models.fake_chat_models import FakeListChatModel


async def chat_model_astream() -> None:
telemetry = init_telemetry("langchain-chat-model-astream")
model = FakeChatModel()
init_telemetry("langchain-chat-model-astream")
model = FakeListChatModel(responses=["Asynchronous streaming chat output."])
chunks = []
async for chunk in model.astream(
"Stream asynchronously.",
config=tracing_config("chat_model_astream"),
):
chunks.append(message_text(chunk))
print("".join(chunks))


if __name__ == "__main__":
asyncio.run(chat_model_astream())
9 changes: 5 additions & 4 deletions python/tracing/langchain/09_chat_model_astream_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,13 @@

import asyncio

from langchain_core.language_models.fake_chat_models import FakeChatModel

from _shared import init_telemetry, tracing_config
from langchain_core.language_models.fake_chat_models import FakeListChatModel


async def chat_model_astream_events() -> None:
telemetry = init_telemetry("langchain-chat-model-astream-events")
model = FakeChatModel()
init_telemetry("langchain-chat-model-astream-events")
model = FakeListChatModel(responses=["Semantic stream events."])
events = []
async for event in model.astream_events(
"Emit semantic stream events.",
Expand All @@ -18,5 +17,7 @@ async def chat_model_astream_events() -> None:
):
events.append(event["event"])
print(events)


if __name__ == "__main__":
asyncio.run(chat_model_astream_events())
44 changes: 41 additions & 3 deletions python/tracing/langchain/11_llm_stream.py
Original file line number Diff line number Diff line change
@@ -1,19 +1,57 @@
"""Legacy string LLM stream."""

from langchain_core.language_models.fake import FakeStreamingListLLM
from collections.abc import Iterator
from typing import Any

from _shared import init_telemetry, tracing_config
from langchain_core.callbacks import CallbackManagerForLLMRun
from langchain_core.language_models.llms import LLM
from langchain_core.outputs import GenerationChunk


class CallbackStreamingLLM(LLM):
"""Deterministic LLM that exercises LangChain's token callback contract."""

response: str

@property
def _llm_type(self) -> str:
return "callback-streaming-list"

def _call(
self,
prompt: str,
stop: list[str] | None = None,
run_manager: CallbackManagerForLLMRun | None = None,
**kwargs: Any,
) -> str:
return self.response

def _stream(
self,
prompt: str,
stop: list[str] | None = None,
run_manager: CallbackManagerForLLMRun | None = None,
**kwargs: Any,
) -> Iterator[GenerationChunk]:
for token in self.response.splitlines(keepends=True):
chunk = GenerationChunk(text=token)
if run_manager is not None:
run_manager.on_llm_new_token(token, chunk=chunk)
yield chunk


def llm_stream() -> None:
telemetry = init_telemetry("langchain-llm-stream")
llm = FakeStreamingListLLM(responses=["tokenized completion"])
init_telemetry("langchain-llm-stream")
llm = CallbackStreamingLLM(response="tokenized completion\nwith callbacks")
chunks = list(
llm.stream(
"Stream this completion.",
config=tracing_config("llm_stream"),
)
)
print("".join(chunks))


if __name__ == "__main__":
llm_stream()
2 changes: 2 additions & 0 deletions python/tracing/langchain/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ Run one example:
python 00_quickstart.py
```

Run the complete bounded set with `python run_all_examples.py`.

## Examples

| Script | LangChain function or behavior |
Expand Down
53 changes: 46 additions & 7 deletions python/tracing/langchain/_shared.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,21 +2,41 @@

from __future__ import annotations

import atexit
import os
from pathlib import Path
from typing import Any

from dotenv import find_dotenv, load_dotenv
from dotenv import load_dotenv
from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel
from langchain_core.messages import AIMessage
from langchain_core.tools import tool
from respan_instrumentation_langchain import add_respan_callback
from respan_tracing import RespanTelemetry

load_dotenv(find_dotenv(), override=False)
ROOT_DIR = Path(__file__).resolve().parents[3]
load_dotenv(ROOT_DIR / ".env", override=True)

RUN_ID = os.getenv("RESPAN_EXAMPLE_RUN_ID", "").strip() or "langchain-local"
_ACTIVE_TELEMETRY: list[RespanTelemetry] = []


class NoopTelemetry:
pass
def flush(self) -> None:
return None


def _flush_telemetry() -> None:
"""Flush every example telemetry instance on normal and exceptional exits."""
while _ACTIVE_TELEMETRY:
telemetry = _ACTIVE_TELEMETRY.pop()
try:
telemetry.flush()
except Exception: # noqa: BLE001,S110 - process-exit flush is best-effort
pass


atexit.register(_flush_telemetry)


def init_telemetry(app_name: str) -> RespanTelemetry | NoopTelemetry:
Expand All @@ -25,22 +45,36 @@ def init_telemetry(app_name: str) -> RespanTelemetry | NoopTelemetry:
if not api_key:
return NoopTelemetry()

return RespanTelemetry(
telemetry = RespanTelemetry(
app_name=app_name,
api_key=api_key,
base_url=os.getenv("RESPAN_BASE_URL", "https://api.respan.ai/api"),
is_auto_instrument=False,
is_batching_enabled=False,
is_enabled=True,
)
_ACTIVE_TELEMETRY.append(telemetry)
return telemetry


def tracing_config(name: str, metadata: dict[str, Any] | None = None) -> dict[str, Any]:
return add_respan_callback(
{
"run_name": name,
"tags": ["respan-langchain-example", name],
"metadata": {"example": name, **(metadata or {})},
"metadata": {
"example": name,
**(metadata or {}),
"respan_params": {
"trace_group_identifier": f"langchain_{name}.workflow",
"custom_identifier": f"{RUN_ID}:{name}",
"metadata": {
"example": "langchain",
"example_run_id": RUN_ID,
"workflow_name": f"langchain_{name}.workflow",
},
},
},
}
)

Expand All @@ -66,7 +100,7 @@ def bind_tools(
self,
tools: Any,
**kwargs: Any,
) -> "ToolCallingFakeMessagesListChatModel":
) -> ToolCallingFakeMessagesListChatModel:
return self


Expand Down Expand Up @@ -107,7 +141,12 @@ def make_openai_chat_model(model_name: str = "gpt-4o-mini") -> Any | None:
"api_key": api_key,
"temperature": 0,
}
base_url = os.getenv("OPENAI_BASE_URL") or os.getenv("RESPAN_OPENAI_BASE_URL")
base_url = (
os.getenv("OPENAI_BASE_URL")
or os.getenv("RESPAN_OPENAI_BASE_URL")
or os.getenv("RESPAN_GATEWAY_BASE_URL")
or os.getenv("RESPAN_BASE_URL")
)
if base_url:
kwargs["base_url"] = base_url
return ChatOpenAI(**kwargs)
Loading