diff --git a/packages/opentelemetry-instrumentation-redis/README.md b/packages/opentelemetry-instrumentation-redis/README.md new file mode 100644 index 0000000000..e69de29bb2 diff --git a/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/__init__.py b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/__init__.py new file mode 100644 index 0000000000..b55a71cef2 --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/__init__.py @@ -0,0 +1,82 @@ +"""OpenTelemetry Redis instrumentation""" + +import logging +import redis +import redis.commands.search +from typing import Collection +from wrapt import wrap_function_wrapper + +from opentelemetry.instrumentation.instrumentor import BaseInstrumentor +from opentelemetry.instrumentation.utils import unwrap + +from opentelemetry.trace import get_tracer + +from opentelemetry.instrumentation.redis.config import Config +from opentelemetry.instrumentation.redis.wrapper import _wrap +from opentelemetry.instrumentation.redis.version import __version__ + +logger = logging.getLogger(__name__) + +_instruments = ("redis >= 4.0.0",) + +WRAPPED_METHODS = [ + { + "package": redis, + "object": "Redis", + "method": "ping", + "span_name": "redis.ping" + }, + { + "package": redis, + "object": "Redis", + "method": "get_connection_kwargs", + "span_name": "redis.getconnectionkwargs" + }, + { + "package": redis.commands.search, + "object": "Search", + "method": "create_index", + "span_name": "create_index", + }, + { + "package": redis.commands.search, + "object": "Search", + "method": "search", + "span_name": "search", + }, +] + +class RedisInstrumentor(BaseInstrumentor): + """An instrumentor for Redis client library.""" + + def __init__(self, exception_logger=None): + print("Redis instrumentor initialized __init__") + super().__init__() + Config.exception_logger = exception_logger + + def instrumentation_dependencies(self) -> Collection[str]: + return _instruments + + def _instrument(self, **kwargs): + print("In redis _instrument") + tracer_provider = kwargs.get("tracer_provider") + tracer = get_tracer(__name__, __version__, tracer_provider) + for wrapped_method in WRAPPED_METHODS: + wrap_package = wrapped_method.get("package") + wrap_object = wrapped_method.get("object") + wrap_method = wrapped_method.get("method") + if getattr(wrap_package, wrap_object, None): + wrap_function_wrapper( + wrap_package, + f"{wrap_object}.{wrap_method}", + _wrap(tracer, wrapped_method), + ) + + def _uninstrument(self, **kwargs): + print("In redis _uninstrument") + for wrapped_method in WRAPPED_METHODS: + wrap_package = wrapped_method.get("package") + wrap_object = wrapped_method.get("object") + wrapped = getattr(wrap_package, wrap_object, None) + if wrapped: + unwrap(wrapped, wrapped_method.get("method")) \ No newline at end of file diff --git a/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/config.py b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/config.py new file mode 100644 index 0000000000..7f991c3ba5 --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/config.py @@ -0,0 +1,2 @@ +class Config: + exception_logger = None \ No newline at end of file diff --git a/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/utils.py b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/utils.py new file mode 100644 index 0000000000..7820408866 --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/utils.py @@ -0,0 +1,23 @@ +import logging +from opentelemetry.instrumentation.redis.config import Config + + +def dont_throw(func): + """ + A decorator that wraps the passed in function and logs exceptions instead of throwing them. + + @param func: The function to wrap + @return: The wrapper function + """ + # Obtain a logger specific to the function's module + logger = logging.getLogger(func.__module__) + + def wrapper(*args, **kwargs): + try: + return func(*args, **kwargs) + except Exception as e: + logger.warning("Failed to execute %s, error: %s", func.__name__, str(e)) + if Config.exception_logger: + Config.exception_logger(e) + + return wrapper diff --git a/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/version.py b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/version.py new file mode 100644 index 0000000000..3aa0d7b3c4 --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/version.py @@ -0,0 +1 @@ +__version__ = "0.0.0" \ No newline at end of file diff --git a/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/wrapper.py b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/wrapper.py new file mode 100644 index 0000000000..4e2fd6ccc6 --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/opentelemetry/instrumentation/redis/wrapper.py @@ -0,0 +1,45 @@ +from opentelemetry import context as context_api +from opentelemetry.semconv.trace import SpanAttributes +from opentelemetry.instrumentation.utils import ( + _SUPPRESS_INSTRUMENTATION_KEY, +) + +def _with_tracer_wrapper(func): + """Helper for providing tracer for wrapper functions.""" + + def _with_tracer(tracer, to_wrap): + def wrapper(wrapped, instance, args, kwargs): + return func(tracer, to_wrap, wrapped, instance, args, kwargs) + + return wrapper + + return _with_tracer + +def _set_span_attribute(span, name, value): + if value is not None: + if value != "": + span.set_attribute(name, value) + return + +def count_or_none(obj): + if obj: + return len(obj) + + return None + +@_with_tracer_wrapper +def _wrap(tracer, to_wrap, wrapped, instance, args, kwargs): + print("In redis _wrap") + """Instruments and calls every function defined in TO_WRAP.""" + if context_api.get_value(_SUPPRESS_INSTRUMENTATION_KEY): + return wrapped(*args, **kwargs) + + name = to_wrap.get("span_name") + with tracer.start_as_current_span(name) as span: + print(f"Name of the span: {name}") + span.set_attribute(SpanAttributes.DB_SYSTEM, "redis") + span.set_attribute(SpanAttributes.DB_OPERATION, to_wrap.get("method")) + _set_span_attribute(span, "db.redis.search.test", 1) + return_value = wrapped(*args, **kwargs) + print(f"Return value: {return_value}") + return return_value \ No newline at end of file diff --git a/packages/opentelemetry-instrumentation-redis/pyproject.toml b/packages/opentelemetry-instrumentation-redis/pyproject.toml new file mode 100644 index 0000000000..c13c98560c --- /dev/null +++ b/packages/opentelemetry-instrumentation-redis/pyproject.toml @@ -0,0 +1,47 @@ +[tool.coverage.run] +branch = true +source = [ "opentelemetry/instrumentation/redis" ] + +[tool.coverage.report] +exclude_lines = [ "if TYPE_CHECKING:" ] +show_missing = true + +[tool.poetry] +name = "opentelemetry-instrumentation-redis" +version = "0.16.6" +description = "OpenTelemetry Redis instrumentation" +authors = [ +] +repository = "https://github.com/traceloop/openllmetry/tree/main/packages/opentelemetry-instrumentation-redis" +license = "Apache-2.0" +readme = "README.md" + +[[tool.poetry.packages]] +include = "opentelemetry/instrumentation/redis" + +[tool.poetry.dependencies] +python = ">=3.9,<4" +opentelemetry-api = "^1.24.0" +opentelemetry-instrumentation = "^0.45b0" +opentelemetry-semantic-conventions = "^0.45b0" +opentelemetry-semantic-conventions-ai = "0.1.1" + +[tool.poetry.group.dev.dependencies] +autopep8 = "2.1.0" +flake8 = "7.0.0" + +[tool.poetry.group.test.dependencies] +redis = "^4.0.0" +pytest = "8.1.1" +pytest-sugar = "1.0.0" +opentelemetry-sdk = "^1.23.0" + +[build-system] +requires = [ "poetry-core" ] +build-backend = "poetry.core.masonry.api" + +[tool.poetry.extras] +instruments = ["redis"] + +[tool.poetry.plugins."opentelemetry_instrumentor"] +redis_client = "opentelemetry.instrumentation.redis:RedisInstrumentor" diff --git a/packages/sample-app/sample_app/redis_app.py b/packages/sample-app/sample_app/redis_app.py new file mode 100644 index 0000000000..92a0169de4 --- /dev/null +++ b/packages/sample-app/sample_app/redis_app.py @@ -0,0 +1,229 @@ +"""_summary_ + Code was taken from https://redis.io/docs/latest/develop/get-started/vector-database/ + It requires a Redis DB running with the RedisSearch and RedisJSON modules enabled. +""" + + +import json +import time + +import numpy as np +import pandas as pd +import redis +import requests +from redis.commands.search.field import ( + NumericField, + TagField, + TextField, + VectorField, +) +from redis.commands.search.indexDefinition import IndexDefinition, IndexType +from redis.commands.search.query import Query +from sentence_transformers import SentenceTransformer + +from traceloop.sdk import Traceloop +from traceloop.sdk.decorators import workflow +from opentelemetry.sdk.trace.export import ConsoleSpanExporter + +#Traceloop.init(exporter=ConsoleSpanExporter()) +Traceloop.init(exporter=ConsoleSpanExporter(), app_name="redis", disable_batch=True) + +url = "https://raw.githubusercontent.com/bsbodden/redis_vss_getting_started/main/data/bikes.json" +response = requests.get(url) +bikes = response.json() + +json.dumps(bikes[0], indent=2) + +client = redis.Redis(host="localhost", port=6379, decode_responses=True) + +res = client.ping() +# >>> True + +pipeline = client.pipeline() +for i, bike in enumerate(bikes, start=1): + redis_key = f"bikes:{i:03}" + pipeline.json().set(redis_key, "$", bike) +res = pipeline.execute() +# >>> [True, True, True, True, True, True, True, True, True, True, True] + +# res = client.json().get("bikes:010", "$.model") +# # >>> ['Summit'] + +# keys = sorted(client.keys("bikes:*")) +# # >>> ['bikes:001', 'bikes:002', ..., 'bikes:011'] + +# descriptions = client.json().mget(keys, "$.description") +# descriptions = [item for sublist in descriptions for item in sublist] +# embedder = SentenceTransformer("msmarco-distilbert-base-v4") +# embeddings = embedder.encode(descriptions).astype(np.float32).tolist() +# VECTOR_DIMENSION = len(embeddings[0]) +# # >>> 768 + +# pipeline = client.pipeline() +# for key, embedding in zip(keys, embeddings): +# pipeline.json().set(key, "$.description_embeddings", embedding) +# pipeline.execute() +# # >>> [True, True, True, True, True, True, True, True, True, True, True] + +# res = client.json().get("bikes:010") +# # >>> +# # { +# # "model": "Summit", +# # "brand": "nHill", +# # "price": 1200, +# # "type": "Mountain Bike", +# # "specs": { +# # "material": "alloy", +# # "weight": "11.3" +# # }, +# # "description": "This budget mountain bike from nHill performs well..." +# # "description_embeddings": [ +# # -0.538114607334137, +# # -0.49465855956077576, +# # -0.025176964700222015, +# # ... +# # ] +# # } + + +# schema = ( +# TextField("$.model", no_stem=True, as_name="model"), +# TextField("$.brand", no_stem=True, as_name="brand"), +# NumericField("$.price", as_name="price"), +# TagField("$.type", as_name="type"), +# TextField("$.description", as_name="description"), +# VectorField( +# "$.description_embeddings", +# "FLAT", +# { +# "TYPE": "FLOAT32", +# "DIM": VECTOR_DIMENSION, +# "DISTANCE_METRIC": "COSINE", +# }, +# as_name="vector", +# ), +# ) +# definition = IndexDefinition(prefix=["bikes:"], index_type=IndexType.JSON) +# res = client.ft("idx:bikes_vss").create_index( +# fields=schema, definition=definition +# ) +# # >>> 'OK' + +# info = client.ft("idx:bikes_vss").info() +# num_docs = info["num_docs"] +# indexing_failures = info["hash_indexing_failures"] +# # print(f"{num_docs} documents indexed with {indexing_failures} failures") +# # >>> 11 documents indexed with 0 failures + +# query = Query("@brand:Peaknetic") +# res = client.ft("idx:bikes_vss").search(query).docs +# # print(res) +# # >>> [Document {'id': 'bikes:008', 'payload': None, 'brand': 'Peaknetic', 'model': 'Soothe Electric bike', 'price': '1950', 'description_embeddings': ... + +# query = Query("@brand:Peaknetic").return_fields("id", "brand", "model", "price") +# res = client.ft("idx:bikes_vss").search(query).docs +# # print(res) +# # >>> [Document {'id': 'bikes:008', 'payload': None, 'brand': 'Peaknetic', 'model': 'Soothe Electric bike', 'price': '1950'}, Document {'id': 'bikes:009', 'payload': None, 'brand': 'Peaknetic', 'model': 'Secto', 'price': '430'}] + +# query = Query("@brand:Peaknetic @price:[0 1000]").return_fields( +# "id", "brand", "model", "price" +# ) +# res = client.ft("idx:bikes_vss").search(query).docs +# # print(res) +# # >>> [Document {'id': 'bikes:009', 'payload': None, 'brand': 'Peaknetic', 'model': 'Secto', 'price': '430'}] + +# queries = [ +# "Bike for small kids", +# "Best Mountain bikes for kids", +# "Cheap Mountain bike for kids", +# "Female specific mountain bike", +# "Road bike for beginners", +# "Commuter bike for people over 60", +# "Comfortable commuter bike", +# "Good bike for college students", +# "Mountain bike for beginners", +# "Vintage bike", +# "Comfortable city bike", +# ] + +# encoded_queries = embedder.encode(queries) +# len(encoded_queries) +# # >>> 11 + +# @workflow("create_query_table") +# def create_query_table(query, queries, encoded_queries, extra_params={}): +# results_list = [] +# for i, encoded_query in enumerate(encoded_queries): +# result_docs = ( +# client.ft("idx:bikes_vss") +# .search( +# query, +# { +# "query_vector": np.array( +# encoded_query, dtype=np.float32 +# ).tobytes() +# } +# | extra_params, +# ) +# .docs +# ) +# for doc in result_docs: +# vector_score = round(1 - float(doc.vector_score), 2) +# results_list.append( +# { +# "query": queries[i], +# "score": vector_score, +# "id": doc.id, +# "brand": doc.brand, +# "model": doc.model, +# "description": doc.description, +# } +# ) + +# # Optional: convert the table to Markdown using Pandas +# queries_table = pd.DataFrame(results_list) +# queries_table.sort_values( +# by=["query", "score"], ascending=[True, False], inplace=True +# ) +# queries_table["query"] = queries_table.groupby("query")["query"].transform( +# lambda x: [x.iloc[0]] + [""] * (len(x) - 1) +# ) +# queries_table["description"] = queries_table["description"].apply( +# lambda x: (x[:497] + "...") if len(x) > 500 else x +# ) +# queries_table.to_markdown(index=False) + + + +# query = ( +# Query("(*)=>[KNN 3 @vector $query_vector AS vector_score]") +# .sort_by("vector_score") +# .return_fields("vector_score", "id", "brand", "model", "description") +# .dialect(2) +# ) + +# create_query_table(query, queries, encoded_queries) +# # >>> | Best Mountain bikes for kids | 0.54 | bikes:003... (+ 32 more results) + +# hybrid_query = ( +# Query("(@brand:Peaknetic)=>[KNN 3 @vector $query_vector AS vector_score]") +# .sort_by("vector_score") +# .return_fields("vector_score", "id", "brand", "model", "description") +# .dialect(2) +# ) +# create_query_table(hybrid_query, queries, encoded_queries) +# # >>> | Best Mountain bikes for kids | 0.3 | bikes:008... (+22 more results) + +# range_query = ( +# Query( +# "@vector:[VECTOR_RANGE $range $query_vector]=>{$YIELD_DISTANCE_AS: vector_score}" +# ) +# .sort_by("vector_score") +# .return_fields("vector_score", "id", "brand", "model", "description") +# .paging(0, 4) +# .dialect(2) +# ) +# create_query_table( +# range_query, queries[:1], encoded_queries[:1], {"range": 0.55} +# ) +# # >>> | Bike for small kids | 0.52 | bikes:001 | Velorim |... (+1 more result) diff --git a/packages/traceloop-sdk/pyproject.toml b/packages/traceloop-sdk/pyproject.toml index 19e1f770df..352ae61a5a 100644 --- a/packages/traceloop-sdk/pyproject.toml +++ b/packages/traceloop-sdk/pyproject.toml @@ -38,6 +38,7 @@ opentelemetry-instrumentation-anthropic = {path="../opentelemetry-instrumentatio opentelemetry-instrumentation-cohere = {path="../opentelemetry-instrumentation-cohere", develop=true} opentelemetry-instrumentation-pinecone = {path="../opentelemetry-instrumentation-pinecone", develop=true} opentelemetry-instrumentation-qdrant = {path="../opentelemetry-instrumentation-qdrant", develop=true} +opentelemetry-instrumentation-redis = {path="../opentelemetry-instrumentation-redis", develop=true} opentelemetry-instrumentation-langchain = {path="../opentelemetry-instrumentation-langchain", develop=true} opentelemetry-instrumentation-chromadb = {path="../opentelemetry-instrumentation-chromadb", develop=true} opentelemetry-instrumentation-transformers = {path="../opentelemetry-instrumentation-transformers", develop=true} diff --git a/packages/traceloop-sdk/traceloop/sdk/instruments.py b/packages/traceloop-sdk/traceloop/sdk/instruments.py index a387ee75c4..073d746c40 100644 --- a/packages/traceloop-sdk/traceloop/sdk/instruments.py +++ b/packages/traceloop-sdk/traceloop/sdk/instruments.py @@ -17,4 +17,5 @@ class Instruments(Enum): REPLICATE = "replicate" VERTEXAI = "vertexai" WATSONX = "watsonx" + REDIS = "redis" WEAVIATE = "weaviate" diff --git a/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py b/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py index 4edc5be98b..ac47a7e475 100644 --- a/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py +++ b/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py @@ -122,6 +122,7 @@ def __new__( init_instrumentations(should_enrich_metrics) instrument_set = True else: + print(Fore.RED + f"Instruments: {instruments}" + Fore.RESET) for instrument in instruments: if instrument == Instruments.OPENAI: if not init_openai_instrumentor(should_enrich_metrics): @@ -182,6 +183,13 @@ def __new__( print(Fore.RESET) else: instrument_set = True + elif instrument == Instruments.REDIS: + if not init_redis_instrumentor(): + print(Fore.RED + "Warning: Redis library does not exist.") + print(Fore.RESET) + else: + print(Fore.RED + "Redis initialized" + Fore.RESET) + instrument_set = True elif instrument == Instruments.REQUESTS: if not init_requests_instrumentor(): print( @@ -422,11 +430,13 @@ def init_tracer_provider(resource: Resource) -> TracerProvider: def init_instrumentations(should_enrich_metrics: bool): + print(Fore.RED + f"init_instrumentations tracing.py" + Fore.RESET) init_openai_instrumentor(should_enrich_metrics) init_anthropic_instrumentor(should_enrich_metrics) init_cohere_instrumentor() init_pinecone_instrumentor() init_qdrant_instrumentor() + init_redis_instrumentor() init_chroma_instrumentor() init_haystack_instrumentor() init_langchain_instrumentor() @@ -509,6 +519,20 @@ def init_qdrant_instrumentor(): instrumentor.instrument() +def init_redis_instrumentor(): + print(Fore.RED + f"Trying to initialize redis instrumentor" + Fore.RESET) + if importlib.util.find_spec("redis") is not None: + print(Fore.RED + "Initializing redis instrumentor" + Fore.RESET) + Telemetry().capture("instrumentation:redis:init") + from opentelemetry.instrumentation.redis import RedisInstrumentor + + instrumentor = RedisInstrumentor( + exception_logger=lambda e: Telemetry().log_exception(e), + ) + if not instrumentor.is_instrumented_by_opentelemetry: + instrumentor.instrument() + + def init_chroma_instrumentor(): if importlib.util.find_spec("chromadb") is not None: Telemetry().capture("instrumentation:chromadb:init") @@ -549,6 +573,7 @@ def init_langchain_instrumentor(): def init_transformers_instrumentor(): + print(Fore.RED + f"Init transformers instrumentor" + Fore.RESET) if importlib.util.find_spec("transformers") is not None: Telemetry().capture("instrumentation:transformers:init") from opentelemetry.instrumentation.transformers import TransformersInstrumentor