mirror of
https://github.com/tiennm99/DocsGPT.git
synced 2026-10-11 03:12:55 +00:00
Keep trace summaries from breaking Logs and fix their counts
A failed trace-summary lookup now leaves the Logs page intact without chips. Tool-call counts include only calls that ran, not their paused, denied or skipped records. A local guardrail that fires unchanged on every streamed segment is recorded once, so it cannot use up the span cap.
This commit is contained in:
1 parent
24796ae61f
commit
9f7f0b2f1e
7 files changed
+134
-19
No files matched your search
@@ -1,10 +1,11 @@
|
||||
"""Analytics and reporting routes."""
|
||||
|
||||
import datetime
|
||||
from typing import Optional, Tuple
|
||||
|
||||
from flask import current_app, jsonify, make_response, request
|
||||
from flask_restx import fields, Namespace, Resource
|
||||
from sqlalchemy import text as _sql_text
|
||||
from sqlalchemy import Connection, text as _sql_text
|
||||
|
||||
from docsgpt.api import api
|
||||
from docsgpt.api.user.base import (
|
||||
@@ -123,7 +124,7 @@ def _trace_branch(name: str, sources_sql: str, scope: str) -> dict:
|
||||
}
|
||||
|
||||
|
||||
def _trace_ref(item: dict):
|
||||
def _trace_ref(item: dict) -> Optional[Tuple[str, str]]:
|
||||
"""The ``(field, value)`` that finds a Logs row's trace, if it has one."""
|
||||
event_type = item["event_type"]
|
||||
row_id = item["id"].split("-", 1)[1]
|
||||
@@ -159,7 +160,9 @@ def _merge_trace_summaries(traces: list) -> dict:
|
||||
}
|
||||
|
||||
|
||||
def _attach_trace_summaries(conn, items: list, *, user_id, agent_id) -> None:
|
||||
def _attach_trace_summaries(
|
||||
conn: Connection, items: list, *, user_id: Optional[str], agent_id: Optional[str]
|
||||
) -> None:
|
||||
"""Add a ``trace`` summary to each Logs row that has a stored trace.
|
||||
|
||||
One batched lookup per link field for the whole page, rather than a join
|
||||
@@ -1282,12 +1285,20 @@ class GetUserLogs(Resource):
|
||||
)
|
||||
results.append(item)
|
||||
if results:
|
||||
with db_readonly() as conn:
|
||||
_attach_trace_summaries(
|
||||
conn,
|
||||
results,
|
||||
user_id=user,
|
||||
agent_id=agent_pg_id if api_key_id else None,
|
||||
# Trace chips are an extra: a failed lookup must leave the
|
||||
# page intact, just without them.
|
||||
try:
|
||||
with db_readonly() as conn:
|
||||
_attach_trace_summaries(
|
||||
conn,
|
||||
results,
|
||||
user_id=user,
|
||||
agent_id=agent_pg_id if api_key_id else None,
|
||||
)
|
||||
except Exception:
|
||||
current_app.logger.warning(
|
||||
"Could not attach trace summaries to the logs page",
|
||||
exc_info=True,
|
||||
)
|
||||
except Exception as err:
|
||||
current_app.logger.error(
|
||||
|
||||
@@ -38,7 +38,23 @@ def _start_guardrail_span(stage: Stage, controls) -> "tracing.Span":
|
||||
)
|
||||
|
||||
|
||||
def _describe_decision(span, decision: StageDecision) -> None:
|
||||
def _content_fired(decision: StageDecision) -> bool:
|
||||
return bool(decision.triggered or decision.blocked or decision.redacted)
|
||||
|
||||
|
||||
def _decision_key(decision: StageDecision) -> tuple:
|
||||
"""Identity of a guardrail outcome, for recording each distinct one once."""
|
||||
return (
|
||||
"guardrail",
|
||||
decision.stage.value,
|
||||
tuple(sorted(v.check for v in decision.triggered)),
|
||||
tuple(sorted(v.check for v in decision.unevaluated)),
|
||||
decision.blocked,
|
||||
decision.redacted,
|
||||
)
|
||||
|
||||
|
||||
def _describe_decision(span: "tracing.Span", decision: StageDecision) -> None:
|
||||
"""Record a stage decision on its span; a firing guardrail drops trace previews."""
|
||||
triggered = [v.check for v in decision.triggered]
|
||||
span.set(
|
||||
@@ -50,7 +66,7 @@ def _describe_decision(span, decision: StageDecision) -> None:
|
||||
"docsgpt.guardrail.unevaluated": [v.check for v in decision.unevaluated] or None,
|
||||
}
|
||||
)
|
||||
if triggered or decision.blocked or decision.redacted:
|
||||
if _content_fired(decision):
|
||||
# The scanned text (or text near it) sits in other spans' previews:
|
||||
# the retrieval query, tool results, the answer. Keep none of it.
|
||||
tracing.mark_content_blocked()
|
||||
@@ -136,10 +152,15 @@ class GuardrailEngine:
|
||||
decision.verdicts = [self._run_control(c, text, stage) for c in controls]
|
||||
self._reduce(decision)
|
||||
# The output guard evaluates every streamed segment; tracing each
|
||||
# clean local scan would bury the trace, so only a firing is kept.
|
||||
# clean local scan would bury the trace, so only a firing is kept,
|
||||
# and a firing that repeats unchanged segment after segment is
|
||||
# recorded once so it cannot use up the trace's span cap.
|
||||
if not decision.clean or decision.unevaluated:
|
||||
with _start_guardrail_span(stage, controls) as span:
|
||||
_describe_decision(span, decision)
|
||||
if tracing.first_occurrence(_decision_key(decision)):
|
||||
with _start_guardrail_span(stage, controls) as span:
|
||||
_describe_decision(span, decision)
|
||||
elif _content_fired(decision):
|
||||
tracing.mark_content_blocked()
|
||||
|
||||
self._record(decision)
|
||||
return decision
|
||||
|
||||
@@ -43,6 +43,7 @@ from docsgpt.tracing.core import (
|
||||
bind,
|
||||
bind_if_unset,
|
||||
current_trace,
|
||||
first_occurrence,
|
||||
mark_content_blocked,
|
||||
span,
|
||||
start_span,
|
||||
@@ -78,6 +79,7 @@ __all__ = [
|
||||
"bind_if_unset",
|
||||
"current_trace",
|
||||
"discard",
|
||||
"first_occurrence",
|
||||
"flush",
|
||||
"mark_content_blocked",
|
||||
"span",
|
||||
|
||||
+34
-5
@@ -272,6 +272,7 @@ class Trace:
|
||||
self.attributes: Dict[str, Any] = {}
|
||||
self.finished = False
|
||||
self.flushed = False
|
||||
self._seen_keys: set = set()
|
||||
self._lock = threading.Lock()
|
||||
self._stacks: Dict[int, List[Any]] = {}
|
||||
|
||||
@@ -288,7 +289,7 @@ class Trace:
|
||||
*,
|
||||
parent: Any = None,
|
||||
attributes: Optional[Dict[str, Any]] = None,
|
||||
):
|
||||
) -> Any:
|
||||
"""Record a new span; returns :data:`NOOP_SPAN` once finished or over the cap."""
|
||||
with self._lock:
|
||||
if self.finished:
|
||||
@@ -344,10 +345,18 @@ class Trace:
|
||||
|
||||
return _undo
|
||||
|
||||
def current_parent(self) -> Optional["_Seed | Span"]:
|
||||
def current_parent(self) -> Optional[Any]:
|
||||
stack = self._stacks.get(threading.get_ident())
|
||||
return stack[-1] if stack else None
|
||||
|
||||
def first_occurrence(self, key: Any) -> bool:
|
||||
"""True the first time ``key`` is seen in this trace, then False."""
|
||||
with self._lock:
|
||||
if key in self._seen_keys:
|
||||
return False
|
||||
self._seen_keys.add(key)
|
||||
return True
|
||||
|
||||
# -- ids --------------------------------------------------------------
|
||||
|
||||
def bind(self, *, only_if_unset: bool = False, **ids: Any) -> None:
|
||||
@@ -409,7 +418,12 @@ class Trace:
|
||||
if s.kind == KIND_RETRIEVAL
|
||||
and not (s.parent_id in by_id and by_id[s.parent_id].kind == KIND_RETRIEVAL)
|
||||
]
|
||||
tools = [s for s in self.spans if s.kind == KIND_TOOL]
|
||||
# Only calls that ran count: a call paused for approval is recorded
|
||||
# again when it runs in the next round, and denied or skipped calls
|
||||
# never ran at all.
|
||||
tools = [
|
||||
s for s in self.spans if s.kind == KIND_TOOL and s.status in (STATUS_OK, STATUS_ERROR)
|
||||
]
|
||||
|
||||
def _tokens(key: str) -> int:
|
||||
total = 0
|
||||
@@ -534,7 +548,9 @@ def activate(trace: Optional[Trace]) -> Iterator[Optional[Trace]]:
|
||||
_current.set(None)
|
||||
|
||||
|
||||
def start_span(kind: str, name: str, *, parent: Any = None, attributes: Optional[Dict[str, Any]] = None):
|
||||
def start_span(
|
||||
kind: str, name: str, *, parent: Any = None, attributes: Optional[Dict[str, Any]] = None
|
||||
) -> Any:
|
||||
"""Start a span in the current trace; returns :data:`NOOP_SPAN` without one."""
|
||||
trace = _current.get()
|
||||
if trace is None:
|
||||
@@ -546,7 +562,9 @@ def start_span(kind: str, name: str, *, parent: Any = None, attributes: Optional
|
||||
return NOOP_SPAN
|
||||
|
||||
|
||||
def span(kind: str, name: str, *, parent: Any = None, attributes: Optional[Dict[str, Any]] = None):
|
||||
def span(
|
||||
kind: str, name: str, *, parent: Any = None, attributes: Optional[Dict[str, Any]] = None
|
||||
) -> Any:
|
||||
"""Context-manager form of :func:`start_span` (errors mark the span and re-raise)."""
|
||||
return start_span(kind, name, parent=parent, attributes=attributes)
|
||||
|
||||
@@ -565,6 +583,17 @@ def bind_if_unset(**ids: Any) -> None:
|
||||
trace.bind(only_if_unset=True, **ids)
|
||||
|
||||
|
||||
def first_occurrence(key: Any) -> bool:
|
||||
"""True the first time ``key`` is seen in the current trace.
|
||||
|
||||
For steps that repeat identically many times in one request (a guardrail
|
||||
re-firing on every streamed segment): record the first, skip the rest, so
|
||||
they cannot use up the span cap. False without an active trace.
|
||||
"""
|
||||
trace = _current.get()
|
||||
return trace.first_occurrence(key) if trace is not None else False
|
||||
|
||||
|
||||
def mark_content_blocked() -> None:
|
||||
"""A guardrail blocked or retracted content: drop every preview from the stored trace."""
|
||||
trace = _current.get()
|
||||
|
||||
@@ -193,3 +193,22 @@ class TestLogsTraceSummaries:
|
||||
|
||||
def test_unknown_event_type_is_rejected(self, app, pg_conn):
|
||||
assert _logs(app, pg_conn, "owner", {"event_type": "bogus"}).status_code == 400
|
||||
|
||||
|
||||
class TestSummaryFailure:
|
||||
def test_page_survives_a_failed_trace_lookup(self, app, pg_conn):
|
||||
from docsgpt.storage.db.repositories.user_logs import UserLogsRepository
|
||||
|
||||
UserLogsRepository(pg_conn).insert(
|
||||
user_id="owner",
|
||||
endpoint="stream_answer",
|
||||
data={"question": "q", "request_id": "req-1"},
|
||||
)
|
||||
with patch(
|
||||
"docsgpt.api.user.analytics.routes.RequestTracesRepository.summaries_for_refs",
|
||||
side_effect=RuntimeError("statement timeout"),
|
||||
):
|
||||
response = _logs(app, pg_conn, "owner", {})
|
||||
assert response.status_code == 200
|
||||
(row,) = response.json["logs"]
|
||||
assert "trace" not in row
|
||||
@@ -388,6 +388,16 @@ class TestTraceSpans:
|
||||
assert span.attributes["docsgpt.guardrail.triggered"] == ["_test_always"]
|
||||
assert self.trace.content_blocked is True
|
||||
|
||||
def test_a_repeating_firing_is_recorded_once(self):
|
||||
"""The output guard re-scans every segment; one firing must not fill the span cap."""
|
||||
engine = GuardrailEngine(
|
||||
_config(controls=[{"check": "_test_always", "stage": "output", "action": "flag"}])
|
||||
)
|
||||
for _ in range(50):
|
||||
engine.evaluate("segment", Stage.OUTPUT)
|
||||
assert len(self.trace.spans) == 1
|
||||
assert self.trace.content_blocked is True
|
||||
|
||||
def test_remote_scan_is_traced_with_judge_nested(self):
|
||||
engine = GuardrailEngine(
|
||||
_config(controls=[{"check": "_test_quick_remote", "stage": "input", "action": "block"}])
|
||||
|
||||
@@ -372,3 +372,26 @@ class TestRecordQuery:
|
||||
tracing.mark_content_blocked()
|
||||
trace.finish()
|
||||
assert "query" not in trace.to_record()["summary"]
|
||||
|
||||
|
||||
class TestToolCount:
|
||||
def test_only_tool_calls_that_ran_are_counted(self):
|
||||
trace = tracing.start_trace(source="stream")
|
||||
with tracing.activate(trace):
|
||||
tracing.start_span(tracing.KIND_TOOL, "paused").end(tracing.STATUS_PENDING)
|
||||
tracing.start_span(tracing.KIND_TOOL, "denied").end(tracing.STATUS_DENIED)
|
||||
tracing.start_span(tracing.KIND_TOOL, "skipped").end(tracing.STATUS_SKIPPED)
|
||||
tracing.start_span(tracing.KIND_TOOL, "ran").end()
|
||||
tracing.start_span(tracing.KIND_TOOL, "failed").end(tracing.STATUS_ERROR)
|
||||
trace.finish()
|
||||
assert trace.summary()["tool_calls"] == 2
|
||||
|
||||
|
||||
class TestFirstOccurrence:
|
||||
def test_key_is_new_once_per_trace(self):
|
||||
trace = tracing.start_trace(source="stream")
|
||||
with tracing.activate(trace):
|
||||
assert tracing.first_occurrence(("a", 1)) is True
|
||||
assert tracing.first_occurrence(("a", 1)) is False
|
||||
assert tracing.first_occurrence(("a", 2)) is True
|
||||
assert tracing.first_occurrence(("a", 3)) is False # no active trace
|
||||
Reference in new issue
Block a user