Skip to content
Merged
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
6 changes: 3 additions & 3 deletions src/fortifyroot/_vendor/VENDOR_MANIFEST.json
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
{
"vendored_at": "2026-06-19T22:19:54.388087",
"vendored_at": "2026-06-20T00:39:46.172213",
"openllmetry_version": "0.52.6",
"git_commit": "f1fbd88483f2",
"git_commit": "01b2d02c2288",
"git_branch": "HEAD",
"git_tag": "fr-v0.52.6.28",
"git_tag": "fr-v0.52.6.29",
"instrumentation_package_policy": {
"opentelemetry-instrumentation-agno": false,
"opentelemetry-instrumentation-alephalpha": false,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,16 @@
"""
from __future__ import annotations

import time
from typing import Any, Generator

from fortifyroot._vendor.opentelemetry.instrumentation.fortifyroot.text_streaming import (
CompletionTextStreamGroup,
)

PROVIDER = "LlamaIndex"
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 LlamaIndexStreamingSafety:
Expand All @@ -26,6 +29,9 @@ class LlamaIndexStreamingSafety:
"""

def __init__(self, span: Any, span_name: str, request_type: str) -> None:
self._span = span
self._stream_started_at = time.perf_counter()
self._first_token_at: float | None = None
self._streams = CompletionTextStreamGroup(
span=span,
provider=PROVIDER,
Expand Down Expand Up @@ -62,6 +68,27 @@ def flush(self, *, segment_index: int = 0) -> str:
"""
return self._streams.flush(key=segment_index) or ""

def record_first_token(self) -> None:
"""Record TTFT once, when the first output-bearing chunk arrives."""
if self._first_token_at is not None:
return
self._first_token_at = time.perf_counter()
_set_streaming_latency_attribute(
self._span,
FR_STREAMING_TIME_TO_FIRST_TOKEN_MS,
self._first_token_at - self._stream_started_at,
)

def record_completion(self) -> None:
"""Record STTG if at least one output-bearing chunk arrived."""
if self._first_token_at is None:
return
_set_streaming_latency_attribute(
self._span,
FR_STREAMING_TIME_TO_GENERATE_MS,
time.perf_counter() - self._first_token_at,
)


# ---------------------------------------------------------------------------
# Internal helpers (best-effort, never raise)
Expand Down Expand Up @@ -91,6 +118,16 @@ def _flush_into_last(response: Any, safety: LlamaIndexStreamingSafety) -> None:
pass


def _set_streaming_latency_attribute(span: Any, name: str, seconds: float) -> None:
if span is not None and span.is_recording() and seconds is not None:
span.set_attribute(name, int(round(max(0, seconds) * 1000)))


def _response_has_output_delta(response: Any) -> bool:
delta = getattr(response, "delta", None)
return isinstance(delta, str) and delta != ""


# ---------------------------------------------------------------------------
# Sync streaming wrapper
# ---------------------------------------------------------------------------
Expand All @@ -106,10 +143,13 @@ def wrap_stream(gen: Generator, safety: LlamaIndexStreamingSafety) -> Generator:
for response in gen:
if pending is not None:
yield pending
if _response_has_output_delta(response):
safety.record_first_token()
_patch_delta(response, safety)
pending = response
if pending is not None:
_flush_into_last(pending, safety)
safety.record_completion()
yield pending


Expand All @@ -132,10 +172,13 @@ async def _gen():
async for response in agen:
if pending is not None:
yield pending
if _response_has_output_delta(response):
safety.record_first_token()
_patch_delta(response, safety)
pending = response
if pending is not None:
_flush_into_last(pending, safety)
safety.record_completion()
yield pending

return _gen()
Loading