From 9b263681aa05a5d85db9ac86aff79c86f565f5b0 Mon Sep 17 00:00:00 2001 From: sayonfortify Date: Fri, 19 Jun 2026 13:13:08 +0530 Subject: [PATCH] chore(vendor): sync OpenLLMetry fr-v0.52.6.25 --- .../_vendor/VENDOR_DEPENDENCIES.json | 6 +- src/fortifyroot/_vendor/VENDOR_MANIFEST.json | 6 +- .../instrumentation/bedrock/__init__.py | 20 ++++-- .../bedrock/streaming_safety.py | 69 ++++++++++++++++++- 4 files changed, 89 insertions(+), 12 deletions(-) diff --git a/src/fortifyroot/_vendor/VENDOR_DEPENDENCIES.json b/src/fortifyroot/_vendor/VENDOR_DEPENDENCIES.json index e904ed4..9654231 100644 --- a/src/fortifyroot/_vendor/VENDOR_DEPENDENCIES.json +++ b/src/fortifyroot/_vendor/VENDOR_DEPENDENCIES.json @@ -52,7 +52,7 @@ }, "test": { "anthropic": ">=0.25.2,<0.83.0", - "boto3": ">=1.35.49,<2", + "boto3": ">=1.34.120,<2", "chromadb": ">=0.5.23,<0.6.0", "google-genai": ">=1.58.0,<1.59.0", "langchain": ">=1.0.0,<2.0.0", @@ -80,10 +80,10 @@ "opentelemetry-sdk": ">=1.38.0,<2", "pandas": ">=1.0.0", "pydantic": "<3", - "pytest": ">=8.3.3,<9", + "pytest": ">=8.2.2,<9", "pytest-asyncio": ">=0.23.7,<1.4.0", "pytest-recording": ">=0.13.1,<0.14.0", - "pytest-sugar": "==1.1.1", + "pytest-sugar": "==1.0.0", "requests": ">=2.31.0,<3", "sqlalchemy": ">=2.0.31,<3", "vcrpy": ">=8.0.0,<9" diff --git a/src/fortifyroot/_vendor/VENDOR_MANIFEST.json b/src/fortifyroot/_vendor/VENDOR_MANIFEST.json index 74b23eb..6f237fe 100644 --- a/src/fortifyroot/_vendor/VENDOR_MANIFEST.json +++ b/src/fortifyroot/_vendor/VENDOR_MANIFEST.json @@ -1,9 +1,9 @@ { - "vendored_at": "2026-06-19T04:00:05.596450", + "vendored_at": "2026-06-19T13:11:58.064294", "openllmetry_version": "0.52.6", - "git_commit": "2f713072095b", + "git_commit": "82ff991d682b", "git_branch": "HEAD", - "git_tag": "fr-v0.52.6.24", + "git_tag": "fr-v0.52.6.25", "instrumentation_package_policy": { "opentelemetry-instrumentation-agno": false, "opentelemetry-instrumentation-alephalpha": false, diff --git a/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/__init__.py b/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/__init__.py index d2b7d6c..05c2d94 100644 --- a/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/__init__.py +++ b/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/__init__.py @@ -263,9 +263,12 @@ def with_instrumentation(*args, **kwargs): # the response stream is fully consumed (which is why # ``start_as_current_span`` cannot be used directly). kwargs = _apply_invoke_prompt_safety(span, kwargs, _BEDROCK_INVOKE_SPAN_NAME) + stream_start_time = time.time() with trace.use_span(span, end_on_exit=False): response = fn(*args, **kwargs) - _handle_stream_call(span, kwargs, response, metric_params, event_logger) + _handle_stream_call( + span, kwargs, response, metric_params, event_logger, stream_start_time + ) return response @@ -303,10 +306,13 @@ def with_instrumentation(*args, **kwargs): kwargs = _apply_converse_prompt_safety(span, kwargs, _BEDROCK_CONVERSE_SPAN_NAME) # ST-10.4: see _instrumented_model_invoke_with_response_stream # for the rationale on use_span(end_on_exit=False). + stream_start_time = time.time() with trace.use_span(span, end_on_exit=False): response = fn(*args, **kwargs) if span.is_recording(): - _handle_converse_stream(span, kwargs, response, metric_params, event_logger) + _handle_converse_stream( + span, kwargs, response, metric_params, event_logger, stream_start_time + ) return response @@ -314,7 +320,9 @@ def with_instrumentation(*args, **kwargs): @dont_throw -def _handle_stream_call(span, kwargs, response, metric_params, event_logger): +def _handle_stream_call( + span, kwargs, response, metric_params, event_logger, stream_start_time +): (provider, model_vendor, model) = _get_vendor_model(kwargs.get("modelId")) request_body = json.loads(kwargs.get("body")) @@ -358,6 +366,7 @@ def stream_done(response_body): StreamingWrapper(response["body"]), span=span, stream_done_callback=stream_done, + stream_start_time=stream_start_time, ) @@ -427,7 +436,9 @@ def _handle_converse(span, kwargs, response, metric_params, event_logger): @dont_throw -def _handle_converse_stream(span, kwargs, response, metric_params, event_logger): +def _handle_converse_stream( + span, kwargs, response, metric_params, event_logger, stream_start_time +): (provider, model_vendor, model) = _get_vendor_model(kwargs.get("modelId")) set_converse_model_span_attributes(span, provider, model, kwargs) @@ -446,6 +457,7 @@ def _handle_converse_stream(span, kwargs, response, metric_params, event_logger) model=model, metric_params=metric_params, event_logger=event_logger, + stream_start_time=stream_start_time, ) diff --git a/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/streaming_safety.py b/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/streaming_safety.py index e3fb516..ac06c2d 100644 --- a/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/streaming_safety.py +++ b/src/fortifyroot/_vendor/opentelemetry/instrumentation/bedrock/streaming_safety.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import time from wrapt import ObjectProxy @@ -20,6 +21,10 @@ from fortifyroot._vendor.opentelemetry.semconv_ai import LLMRequestTypeValues +FR_STREAMING_TIME_TO_FIRST_TOKEN_MS = "fortifyroot.llm.streaming.time_to_first_token_ms" +FR_STREAMING_TIME_TO_GENERATE_MS = "fortifyroot.llm.streaming.time_to_generate_ms" + + class _BedrockChunkStreamingSafety: def __init__(self, *, span, span_name: str, request_type: str): self._streams = CompletionTextStreamGroup( @@ -144,12 +149,17 @@ def _payload_text(self, payload): class BedrockInvokeSafetyStreamingWrapper(ObjectProxy): - def __init__(self, response, *, span, stream_done_callback=None): + def __init__(self, response, *, span, stream_done_callback=None, stream_start_time=None): super().__init__(response) + self._self_span = span self._self_stream_done_callback = stream_done_callback self._self_accumulating_body = {} self._self_pending_event = None + self._self_stream_start_time = ( + time.time() if stream_start_time is None else stream_start_time + ) + self._self_first_token_time = None self._self_streaming_safety = _BedrockChunkStreamingSafety( span=span, span_name=getattr(span, "name", "bedrock.completion"), @@ -173,6 +183,7 @@ def __iter__(self): self._accumulate_event(self._self_pending_event) yield self._self_pending_event + self._finish_streaming_latency() if self._self_stream_done_callback: self._self_stream_done_callback(self._self_accumulating_body) @@ -181,6 +192,8 @@ def _accumulate_event(self, event): if not isinstance(payload, dict): return + self._observe_streaming_text(self._self_streaming_safety._payload_text(payload)) + event_type = payload.get("type") if event_type is None: self._accumulate_events(payload) @@ -217,6 +230,25 @@ def _accumulate_events(self, payload): else: self._self_accumulating_body[key] = value + def _observe_streaming_text(self, text): + if not isinstance(text, str) or not text or self._self_first_token_time is not None: + return + now = time.time() + self._self_first_token_time = now + if self._self_span.is_recording(): + self._self_span.set_attribute( + FR_STREAMING_TIME_TO_FIRST_TOKEN_MS, + round((now - self._self_stream_start_time) * 1000), + ) + + def _finish_streaming_latency(self): + if self._self_first_token_time is None or not self._self_span.is_recording(): + return + self._self_span.set_attribute( + FR_STREAMING_TIME_TO_GENERATE_MS, + round((time.time() - self._self_first_token_time) * 1000), + ) + class _BedrockConverseStreamingSafety: def __init__(self, *, span, span_name: str): @@ -321,6 +353,7 @@ def __init__( model, metric_params, event_logger, + stream_start_time=None, ): super().__init__(response) @@ -333,6 +366,10 @@ def __init__( self._self_role = "unknown" self._self_response_msg = [] self._self_span_ended = False + self._self_stream_start_time = ( + time.time() if stream_start_time is None else stream_start_time + ) + self._self_first_token_time = None self._self_streaming_safety = _BedrockConverseStreamingSafety( span=span, span_name=getattr(span, "name", "bedrock.converse"), @@ -356,6 +393,7 @@ def __iter__(self): yield self._self_pending_event if not self._self_span_ended: + self._finish_streaming_latency() self._self_span.end() self._self_span_ended = True @@ -369,10 +407,12 @@ def _observe_event(self, event): if "contentBlockDelta" in event: delta_text = ((event.get("contentBlockDelta") or {}).get("delta") or {}).get("text") if isinstance(delta_text, str): + self._observe_streaming_text(delta_text) self._self_response_msg.append(delta_text) elif "contentBlockStart" in event: start_text = ((event.get("contentBlockStart") or {}).get("start") or {}).get("text") if isinstance(start_text, str): + self._observe_streaming_text(start_text) self._self_response_msg.append(start_text) if "messageStop" in event: @@ -402,15 +442,38 @@ def _observe_event(self, event): ) converse_usage_record(self._self_span, metadata, self._self_metric_params) if not self._self_span_ended: + self._finish_streaming_latency() self._self_span.end() self._self_span_ended = True + def _observe_streaming_text(self, text): + if not isinstance(text, str) or not text or self._self_first_token_time is not None: + return + now = time.time() + self._self_first_token_time = now + if self._self_span.is_recording(): + self._self_span.set_attribute( + FR_STREAMING_TIME_TO_FIRST_TOKEN_MS, + round((now - self._self_stream_start_time) * 1000), + ) -def create_invoke_stream_wrapper(response, *, span, stream_done_callback=None): + def _finish_streaming_latency(self): + if self._self_first_token_time is None or not self._self_span.is_recording(): + return + self._self_span.set_attribute( + FR_STREAMING_TIME_TO_GENERATE_MS, + round((time.time() - self._self_first_token_time) * 1000), + ) + + +def create_invoke_stream_wrapper( + response, *, span, stream_done_callback=None, stream_start_time=None +): return BedrockInvokeSafetyStreamingWrapper( response, span=span, stream_done_callback=stream_done_callback, + stream_start_time=stream_start_time, ) @@ -422,6 +485,7 @@ def create_converse_stream_wrapper( model, metric_params, event_logger, + stream_start_time=None, ): return BedrockConverseSafetyStream( response, @@ -430,6 +494,7 @@ def create_converse_stream_wrapper( model=model, metric_params=metric_params, event_logger=event_logger, + stream_start_time=stream_start_time, )