From 585abcb97e049a9759c5022f9a9a9a2595cbfc60 Mon Sep 17 00:00:00 2001 From: "my.nguyen" Date: Thu, 24 Sep 2026 17:19:03 +0700 Subject: [PATCH] feat(gooddata-eval): add data obfuscation evaluator Add the agentic_obfuscation kind, which checks that canaries are masked in the stored conversation and the Langfuse trace. jira: QA-29442 risk: low --- packages/gooddata-eval/README.md | 56 +- .../src/gooddata_eval/cli/agentic_runner.py | 19 +- .../src/gooddata_eval/cli/main.py | 13 +- .../core/agentic/_obfuscation_check.py | 223 +++++++ .../core/agentic/_obfuscation_observe.py | 140 +++++ .../core/agentic/_obfuscation_sinks.py | 172 ++++++ .../core/agentic/_trace_linker.py | 5 +- .../gooddata_eval/core/agentic/obfuscation.py | 566 +++++++++++++++++ .../src/gooddata_eval/core/chat/sse_client.py | 13 +- .../core/dataset/langfuse_source.py | 29 +- .../src/gooddata_eval/core/langfuse/client.py | 78 ++- .../core/langfuse/observations.py | 4 +- .../src/gooddata_eval/core/models.py | 15 +- .../tests/test_agentic_obfuscation.py | 581 ++++++++++++++++++ .../tests/test_agentic_runner.py | 1 + packages/gooddata-eval/tests/test_cli.py | 6 + .../tests/test_langfuse_client.py | 93 +++ .../tests/test_langfuse_source.py | 33 + packages/gooddata-eval/tests/test_models.py | 21 + .../tests/test_obfuscation_check.py | 197 ++++++ .../gooddata-eval/tests/test_sse_client.py | 9 + .../gooddata-eval/tests/test_trace_linker.py | 1 + 22 files changed, 2257 insertions(+), 18 deletions(-) create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_check.py create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_observe.py create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_sinks.py create mode 100644 packages/gooddata-eval/src/gooddata_eval/core/agentic/obfuscation.py create mode 100644 packages/gooddata-eval/tests/test_agentic_obfuscation.py create mode 100644 packages/gooddata-eval/tests/test_obfuscation_check.py diff --git a/packages/gooddata-eval/README.md b/packages/gooddata-eval/README.md index ac0b2c72e..5ece80f41 100644 --- a/packages/gooddata-eval/README.md +++ b/packages/gooddata-eval/README.md @@ -136,7 +136,7 @@ gd-eval run \ | `--reasoning-effort LEVEL` | server default | `LOW`, `MEDIUM` or `HIGH`, sent as `options.reasoningEffort` on every chat message. Requires the `enableGenAiReasoningEffort` feature flag on the target organization — without it the server ignores the value. Applies to chat items only; `dashboard_summary` items go through the summary endpoint, which has no such option. | **Concurrency and workspace safety.** Agentic kinds that create workspace objects -(`agentic_metric_skill`, `agentic_alert_skill`, `agentic_conversation`, `agentic_kda_skill`) always run one at a +(`agentic_metric_skill`, `agentic_alert_skill`, `agentic_conversation`, `agentic_kda_skill`, `agentic_obfuscation`) always run one at a time whatever `--concurrency` says — a metric or alert created and dropped mid-run would otherwise be visible to another item reading the same catalog. **That protection is for the agentic kinds only:** the single-turn `metric_skill` and `alert_skill` kinds are still fanned out and the agent performs the same server-side writes on @@ -535,7 +535,8 @@ A dataset is a folder of `.json` files, one per question: ``` Supported `test_kind` values: `visualization`, `metric_skill`, `alert_skill`, -`search_tool`, `general_question`, `guardrail`, `dashboard_summary`. +`search_tool`, `general_question`, `guardrail`, `dashboard_summary`, and the agentic kinds +(`agentic_obfuscation` is described below). ### `dashboard_summary` items @@ -574,6 +575,53 @@ The `expected_output` rubric: Each criterion is scored independently by the LLM judge, so `quality_score` is the fraction of satisfied criteria. +### `agentic_obfuscation` items + +gen-ai masks sensitive values out of the conversation it stores and the trace it exports to +Langfuse, while the model and the user's stream keep what was typed. So these items are not +graded on the answer: they plant synthetic *canaries* in one or more user turns, then read back +the stored conversation (`GET …/chat/conversations/{id}/items`) and every Langfuse trace of the +session, and decide by exact substring. No LLM decides whether a value leaked. + +```json +{ + "id": "obfuscation-001", + "dataset_name": "agent_obfuscation", + "test_kind": "agentic_obfuscation", + "question": ["My email is qa.canary5a1e@example.invalid. Show Total Sales by month.", "Now by quarter instead."], + "expected_output": { + "status": "enforced", + "canaries": [ + {"nonce": "canary5a1e", "value": "qa.canary5a1e@example.invalid", "class": "EMAIL", + "absent_from": ["conversation_db", "langfuse_trace"], "mask_marker_present": "[EMAIL]"} + ] + } +} +``` + +- `question` is a string, or a list of turns sent in order to one conversation (from Langfuse, an + input list). A question that is itself a JSON document must be stored in Langfuse as + `{"query": ""}`: Langfuse parses a JSON-looking string input into an object. +- A canary lists the sinks it must be `absent_from` and the sinks it must stay `present_in`. With + `status: known_limitation` a `present_in` value is a known gap: the item fails once it closes, so + the fixture is flipped deliberately. With `status: enforced` it guards a value that is not + sensitive: masking it fails the item as `OVER-MASKED`. `record_only_paths` names sink paths + reported but not gated. +- `expected_turn_rejected: {"status_code": 422, "reason": "DATA_OBFUSCATION_CONTENT_REJECTED"}` + expects the first turn to be refused. +- `observe` (`automation_match`, `metric_title`, `stream_markers`) reports what the chat created + or streamed, never gates it, and deletes the alert, export or metric it recognises. +- Every turn needs a non-sensitive fragment of 12+ word characters, the *anchor*: a sink read back + without every anchor is reported as blind, never as clean. `anchors` overrides the derived ones. +- Langfuse is polled for the session's traces for 60 s (`GD_EVAL_OBFUSCATION_LANGFUSE_TIMEOUT_SEC`). + When none arrives the run fails with `TRACE_NOT_FOUND`, the report's `trace_found` is false and + `obfuscation_trace_found` is 0 -- a missing trace points at the export, never counts as a pass. + +A leak in any run fails the item whatever `--gate` says. Before the first item of a workspace a +preflight probe checks that both legs mask and both read-backs work, and stops every item of that +workspace when they do not. The kind needs the `LANGFUSE_*` credentials of the Langfuse project +the environment exports to: the traces are read, not only scored. + ## Supported test kinds | test_kind | What the agent must produce | Extra required | @@ -585,6 +633,7 @@ is the fraction of satisfied criteria. | `general_question` | Text answer judged by LLM | `[llm-judge]` | | `guardrail` | Refusal/redirect (visualization response auto-fails) | `[llm-judge]` | | `dashboard_summary` | Dashboard summary (via `/summary` endpoint) scored against a rubric by LLM | `[llm-judge]` | +| `agentic_obfuscation` | Canaries masked in the stored conversation and the Langfuse trace (exact match) | `LANGFUSE_*` | ## Optional extras @@ -625,6 +674,9 @@ the item's own root span. On the agentic path each score is mirrored onto the ag | `quality_score` | Fraction of strict check flags that are `True` (0.0–1.0). Shown in CLI as a percentage. | | `value_score` | Weighted blend: 0.6 × quality + 0.2 × speed (speed = max(0, 1 − latency/60s)). | | `latency_s` | Average per-run latency in seconds. | +| `obfuscation_pass` / `obfuscation_no_leak` | `agentic_obfuscation`: the run's verdict, and whether no canary leaked. The comment names what failed, with canary values replaced by their class. | +| `obfuscation_trace_found` | `agentic_obfuscation`: 0 when Langfuse held no trace for the conversation within the wait (`TRACE_NOT_FOUND`). | +| `obfuscation_record_only_hit` | `agentic_obfuscation`: 1 when a record-only path held a canary; the comment carries every record-only observation. | | `provider_type` | Model vendor + gateway label (e.g. `ANTHROPIC`, `BEDROCK/ANTHROPIC`, `AZURE/OPENAI`). Stored in Langfuse trace metadata and tags. | Score names carry no K; K and the gate are on the dataset-run metadata as `eval_k` and `eval_gate`. diff --git a/packages/gooddata-eval/src/gooddata_eval/cli/agentic_runner.py b/packages/gooddata-eval/src/gooddata_eval/cli/agentic_runner.py index 20885a85c..6c3d35a53 100644 --- a/packages/gooddata-eval/src/gooddata_eval/cli/agentic_runner.py +++ b/packages/gooddata-eval/src/gooddata_eval/cli/agentic_runner.py @@ -18,6 +18,7 @@ from gooddata_eval.core.agentic.guardrail import evaluate_agentic_guardrail from gooddata_eval.core.agentic.kda_skill import evaluate_agentic_kda_skill from gooddata_eval.core.agentic.metric_skill import evaluate_agentic_metric_skill +from gooddata_eval.core.agentic.obfuscation import evaluate_agentic_obfuscation from gooddata_eval.core.agentic.search_tool import evaluate_agentic_search_tool from gooddata_eval.core.agentic.visualization import evaluate_agentic_visualization from gooddata_eval.core.agentic.what_if import evaluate_agentic_what_if @@ -49,6 +50,7 @@ class _LfKw(TypedDict, total=False): "agentic_conversation", "agentic_kda_skill", "agentic_what_if", + "agentic_obfuscation", } ) @@ -88,7 +90,9 @@ class _LfKw(TypedDict, total=False): # metric skill. agentic_kda_skill is here on suspicion rather than proof: it triggers # create_key_driver_analysis with no cleanup, and while the evaluator only ever reads that # call's ARGUMENTS -- never a created object id -- whether the platform persists anything is -# unverified. Move it to the allowlist once someone confirms it does not. +# unverified. Move it to the allowlist once someone confirms it does not. agentic_obfuscation +# items may ask for an alert, a scheduled export or a metric; it deletes what it recognises, +# but that is cleanup, not read-only. # # agentic_dashboard_skill is absent by default rather than by evidence: gen-ai holds the draft and # any chart it authors in conversation state and writes neither until a user saves from the UI, so @@ -278,6 +282,19 @@ def _dispatch_agentic( agent_id=agent_id, **lf_kw, ) + elif kind == "agentic_obfuscation": + return evaluate_agentic_obfuscation( + host=host, + token=token, + workspace_id=workspace_id, + question=item.question, + expected_output=eo if isinstance(eo, dict) else {}, + k=k, + turns=item.turns, + gate=gate, + agent_id=agent_id, + **lf_kw, + ) elif kind == "agentic_conversation": fixture_data = eo.get("fixture") or eo if isinstance(eo, dict) else {} return evaluate_agentic_conversation( diff --git a/packages/gooddata-eval/src/gooddata_eval/cli/main.py b/packages/gooddata-eval/src/gooddata_eval/cli/main.py index a724a8bac..c54ec2e77 100644 --- a/packages/gooddata-eval/src/gooddata_eval/cli/main.py +++ b/packages/gooddata-eval/src/gooddata_eval/cli/main.py @@ -66,6 +66,17 @@ def close(self) -> None: backend.close() +def _positive_int(value: str) -> int: + """An argparse type for counts that must be at least 1.""" + try: + number = int(value) + except ValueError as exc: + raise argparse.ArgumentTypeError(f"not an integer: {value!r}") from exc + if number < 1: + raise argparse.ArgumentTypeError(f"must be at least 1, got {number}") + return number + + def _build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(prog="gd-eval", description="Evaluate the GoodData AI agent.") sub = parser.add_subparsers(dest="command", required=True) @@ -102,7 +113,7 @@ def _build_parser() -> argparse.ArgumentParser: "Default: workspace's current active model." ), ) - run.add_argument("--runs", type=int, default=2, help="Independent runs per item. Default 2.") + run.add_argument("--runs", type=_positive_int, default=2, help="Independent runs per item. Default 2.") run.add_argument( "--gate", choices=get_args(EvalGate), diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_check.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_check.py new file mode 100644 index 000000000..a9407a525 --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_check.py @@ -0,0 +1,223 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +"""Deterministic leak verdict for the agentic_obfuscation kind. + +Pure functions only: the caller collects the sinks, this module decides. A sink is any +JSON-like document -- the stored conversation read back over the API, or every Langfuse +trace of the conversation's session. A canary is found by exact substring over every string +leaf, so the verdict never depends on an LLM. JSON serialised inside a string (Langfuse +keeps span input that way) is decoded and walked too. + +A sink that cannot be seen must not pass an absence check. Every turn therefore carries an +anchor, a non-sensitive fragment of the question, and a sink that does not show every +anchor is reported as blind instead of clean. +""" + +from __future__ import annotations + +import json +import re +from collections.abc import Iterator, Mapping, Sequence +from dataclasses import dataclass, field +from typing import Any + +SINK_CONVERSATION_DB = "conversation_db" +SINK_LANGFUSE_TRACE = "langfuse_trace" +SINKS = (SINK_CONVERSATION_DB, SINK_LANGFUSE_TRACE) + +# Classes whose value may be written with separators the detector ignores, so a leak of +# "4916 3385 0608 2832" must still match the compact nonce "4916338506082832". +_DIGIT_CLASSES = frozenset({"CREDIT_CARD", "IBAN"}) +_SEPARATORS = re.compile(r"[\s\-]") +# Placeholder the AI-on-data trace redaction writes over Langfuse input and output. Behind it +# the Langfuse leg shows nothing, so an absence check there would prove nothing. +AI_ON_DATA_REDACTION = "" +_ANCHOR_RUN = re.compile(r"[^\W_](?:[\w ]*[^\W_])?") +_ANCHOR_MIN = 12 +_ANCHOR_MAX = 40 + + +@dataclass(frozen=True) +class Occurrence: + sink: str + path: str + needle: str + + +@dataclass +class CaseVerdict: + failures: list[str] = field(default_factory=list) + notes: list[str] = field(default_factory=list) + + @property + def passed(self) -> bool: + return not self.failures + + +def iter_strings(value: Any, path: str = "$") -> Iterator[tuple[str, str]]: + """Yield ``(path, text)`` for every string or number leaf and every key, decoding JSON held in a string. + + A string that parses as a JSON document is walked instead of being yielded whole, so a + value spelled with escapes (``\\u0040``) is seen decoded and reported once. Numbers are + leaves too: a numeric secret in a pasted JSON document decodes to an int. Keys are text a + user can type as well -- a pasted config may key on an email -- so each is yielded at the + path of its object with a ```` suffix. + """ + if isinstance(value, bool): + return + if isinstance(value, int | float): + yield path, str(value) + return + if isinstance(value, str): + if value.lstrip()[:1] in ("{", "["): + try: + decoded = json.loads(value) + except ValueError: + pass + else: + yield from iter_strings(decoded, f"{path}") + return + yield path, value + elif isinstance(value, Mapping): + for key, child in value.items(): + if isinstance(key, str): + yield f"{path}.{key}", key + yield from iter_strings(child, f"{path}.{key}") + elif isinstance(value, Sequence) and not isinstance(value, bytes | bytearray): + for index, child in enumerate(value): + yield from iter_strings(child, f"{path}[{index}]") + + +def canary_needles(canary: Mapping[str, Any]) -> list[str]: + """Every spelling of the canary a leak may take: nonce, value and its declared variants.""" + keys = ("nonce", "value", "unescaped_value", "decoded_value", "compact_value") + needles = [canary[key] for key in keys if isinstance(canary.get(key), str) and canary[key]] + return list(dict.fromkeys(needles)) + + +def find_canary(canary: Mapping[str, Any], sink: str, document: Any) -> list[Occurrence]: + needles = canary_needles(canary) + compact = {_SEPARATORS.sub("", n) for n in needles} if canary.get("class") in _DIGIT_CLASSES else set() + found: list[Occurrence] = [] + for path, text in iter_strings(document): + hit = next((n for n in needles if n in text), None) + if hit is None and compact: + squeezed = _SEPARATORS.sub("", text) + hit = next((n for n in compact if n in squeezed), None) + if hit is not None: + found.append(Occurrence(sink, path, hit)) + return found + + +def contains(document: Any, needle: str) -> bool: + return any(needle in text for _, text in iter_strings(document)) + + +def derive_anchor(question: str, canaries: Sequence[Mapping[str, Any]]) -> str | None: + """The longest plain-word run of the question once every canary spelling is cut out. + + Word characters and spaces only, so escaping, JSON quoting and masking next to a canary + cannot alter it between what was sent and what a sink stores. + """ + fragments = [question] + # Longest spelling first: cutting the nonce out of "john\\.canary2d8e@..." first would + # leave "john\\." behind, and the anchor would then end in text the mask replaces. + needles = sorted({n for canary in canaries for n in canary_needles(canary)}, key=len, reverse=True) + for needle in needles: + fragments = [piece for fragment in fragments for piece in fragment.split(needle)] + runs = [match.group(0) for fragment in fragments for match in _ANCHOR_RUN.finditer(fragment)] + best = max(runs, key=len, default="") + if len(best) < _ANCHOR_MIN: + return None + return best[:_ANCHOR_MAX].rstrip() + + +def _declared_sinks(canary: Mapping[str, Any], key: str) -> list[str]: + sinks = canary.get(key) or [] + unknown = [s for s in sinks if s not in SINKS] + if unknown: + raise ValueError(f"canary {canary.get('nonce')!r} names unknown sink(s) {unknown} in {key}") + return list(sinks) + + +def _format(occurrences: Sequence[Occurrence], limit: int = 5) -> str: + shown = ", ".join(f"{o.path} ({o.needle!r})" for o in occurrences[:limit]) + more = len(occurrences) - limit + return shown + (f" and {more} more" if more > 0 else "") + + +def evaluate_case( + expected: Mapping[str, Any], + sinks: Mapping[str, Any], + anchors: Sequence[str], + *, + turn_rejected: bool = False, +) -> CaseVerdict: + """Decide one item from the collected sinks. + + ``sinks`` maps a sink name to its document, or to ``None`` when it was not collected + (only allowed for the Langfuse leg of a rejected turn, which may export nothing). + """ + verdict = CaseVerdict() + status = expected.get("status", "enforced") + visible: dict[str, bool] = {} + + for sink in SINKS: + document = sinks.get(sink) + if document is None: + visible[sink] = False + if not turn_rejected: + verdict.failures.append(f"{sink}: not collected, so no absence claim can be made") + continue + if sink == SINK_LANGFUSE_TRACE and contains(document, AI_ON_DATA_REDACTION): + verdict.failures.append( + f"{sink}: trace input/output is replaced by {AI_ON_DATA_REDACTION!r} (enableAiOnData with " + "enableGenAiTraceRedaction), so the obfuscation leg cannot be observed here" + ) + visible[sink] = False + continue + missing = [anchor for anchor in anchors if not contains(document, anchor)] + if missing and not turn_rejected: + verdict.failures.append(f"{sink}: blind -- anchor(s) {missing} not found, the read-back is incomplete") + visible[sink] = False + continue + visible[sink] = True + + for canary in expected.get("canaries", []): + label = f"{canary.get('class')} {canary.get('nonce')!r}" + for sink in _declared_sinks(canary, "absent_from"): + # A blind sink still convicts: a canary it does show is a leak all the same. + if sinks.get(sink) is None: + continue + occurrences = find_canary(canary, sink, sinks[sink]) + # Paths the item declares as not yet decided (e.g. conversation state the SC does + # not rule on) are reported, never gated, until a decision turns them into leaks. + record_only = [re.compile(p) for p in canary.get("record_only_paths") or []] + recorded = [o for o in occurrences if any(r.search(o.path) for r in record_only)] + gated = [o for o in occurrences if o not in recorded] + if gated: + verdict.failures.append(f"LEAK {label} in {sink}: {_format(gated)}") + if recorded: + verdict.notes.append(f"RECORDED {label} in {sink} (record-only path, not gated): {_format(recorded)}") + marker = canary.get("mask_marker_present") + if marker and not turn_rejected and visible.get(sink) and not contains(sinks[sink], marker): + verdict.failures.append(f"{label}: mask marker {marker!r} missing from {sink}") + for sink in _declared_sinks(canary, "present_in"): + if not visible.get(sink): + continue + if not find_canary(canary, sink, sinks[sink]): + if status == "known_limitation": + flip = (expected.get("known_limitation") or {}).get("flip_when", "") + verdict.failures.append( + f"{label} expected in {sink} ({status}) but is now masked -- the limitation no longer " + f"reproduces, update the fixture deliberately. flip_when: {flip}" + ) + else: + # An enforced present_in guards a value that is not sensitive: masking it is + # the defect (a false positive), not progress. + verdict.failures.append( + f"OVER-MASKED {label} in {sink}: expected unchanged but it was masked (false positive)" + ) + + if turn_rejected and sinks.get(SINK_LANGFUSE_TRACE) is None: + verdict.notes.append("rejected turn exported no Langfuse trace; only the database leg was asserted") + return verdict diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_observe.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_observe.py new file mode 100644 index 000000000..d303e92af --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_observe.py @@ -0,0 +1,140 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +"""Record-only observations for the agentic_obfuscation kind's multi-turn items. + +Some multi-turn behaviour has no agreed expectation yet -- whether an alert still reaches the +email the user typed once history is masked, or a filter still holds its value on the next +turn. An item lists what to look at under ``expected_output.observe``; this module reports +it and never decides a verdict. It also deletes what the chat created (an alert, a scheduled +export, a metric), so a persistent workspace can run the dataset again and again. + +``observe`` keys: + automation_match substring identifying the automation this item makes (a metric id, a + dashboard id); the recipients it ended up with are reported + metric_title title of the metric this item makes; its MAQL is reported + stream_markers strings to look for in the last turn's streamed text and visualizations +""" + +from __future__ import annotations + +import json +from collections.abc import Callable +from typing import Any + +import httpx + +_ENTITIES = "application/vnd.gooddata.api+json" +_PAGE_SIZE = 500 +_MAX_PAGES = 20 + + +class WorkspaceEntities: + """Just enough of the workspace entities API to list and delete what a chat created.""" + + def __init__(self, http: httpx.Client, host: str, workspace_id: str) -> None: + self._http = http + self._base = f"{host.rstrip('/')}/api/v1/entities/workspaces/{workspace_id}" + + def list(self, entity: str) -> dict[str, dict[str, Any]]: + """Every entity of the kind, all pages: a created one missed here is never cleaned up.""" + found: dict[str, dict[str, Any]] = {} + for page in range(_MAX_PAGES): + resp = self._http.get( + f"{self._base}/{entity}", params={"size": _PAGE_SIZE, "page": page}, headers={"Accept": _ENTITIES} + ) + resp.raise_for_status() + data = resp.json().get("data", []) + found.update((item["id"], item) for item in data) + if len(data) < _PAGE_SIZE: + return found + raise httpx.HTTPError(f"{entity}: more than {_MAX_PAGES * _PAGE_SIZE} entities, listing stopped") + + def delete(self, entity: str, entity_id: str) -> None: + resp = self._http.delete(f"{self._base}/{entity}/{entity_id}", headers={"Accept": _ENTITIES}) + # Already gone -- the agent, or the kind's own cleanup, removed it first. + if resp.status_code != 404: + resp.raise_for_status() + + +def snapshot(entities: WorkspaceEntities, spec: dict[str, Any]) -> dict[str, set[str]]: + """Ids that existed before the item ran, for the entities its ``observe`` spec watches.""" + taken: dict[str, set[str]] = {} + if spec.get("automation_match"): + taken["automations"] = set(entities.list("automations")) + if spec.get("metric_title"): + taken["metrics"] = set(entities.list("metrics")) + return taken + + +def _recipients(automation: dict[str, Any]) -> str: + attributes = automation.get("attributes") or {} + external = [r.get("email") for r in attributes.get("externalRecipients") or []] + internal = [r.get("id") for r in ((automation.get("relationships") or {}).get("recipients") or {}).get("data", [])] + return f"external={external} internal={internal}" + + +def observe_and_clean( + entities: WorkspaceEntities, spec: dict[str, Any], before: dict[str, set[str]], last_stream: str +) -> list[str]: + """Report what the item's ``observe`` spec asks about, then delete what the chat created. + + Runs from a ``finally``: a failed list or delete becomes a note, so it can neither replace + the item's verdict nor stop the rest of the cleanup. + """ + notes: list[str] = [] + _observe_step("automations", _observe_automations, entities, spec, before, notes) + _observe_step("metrics", _observe_metrics, entities, spec, before, notes) + notes.extend( + f"OBSERVED last turn stream {'contains' if marker in last_stream else 'lacks'} {marker!r}" + for marker in spec.get("stream_markers") or [] + ) + return notes + + +def _observe_step( + name: str, + step: Callable[[WorkspaceEntities, dict[str, Any], dict[str, set[str]], list[str]], None], + entities: WorkspaceEntities, + spec: dict[str, Any], + before: dict[str, set[str]], + notes: list[str], +) -> None: + try: + step(entities, spec, before, notes) + except httpx.HTTPError as exc: + notes.append(f"OBSERVED cleanup failed ({name}): {exc}") + + +def _delete(entities: WorkspaceEntities, entity: str, entity_id: str, notes: list[str]) -> None: + try: + entities.delete(entity, entity_id) + except httpx.HTTPError as exc: + notes.append(f"OBSERVED could not delete {entity} {entity_id}: {exc}") + + +def _observe_automations( + entities: WorkspaceEntities, spec: dict[str, Any], before: dict[str, set[str]], notes: list[str] +) -> None: + match = spec.get("automation_match") + if match: + created = {i: a for i, a in entities.list("automations").items() if i not in before["automations"]} + mine = {i: a for i, a in created.items() if match in json.dumps(a)} + if not mine: + notes.append(f"OBSERVED no automation matching {match!r} was created") + for automation_id, automation in mine.items(): + notes.append(f"OBSERVED automation {automation_id} created with {_recipients(automation)}") + _delete(entities, "automations", automation_id, notes) + + +def _observe_metrics( + entities: WorkspaceEntities, spec: dict[str, Any], before: dict[str, set[str]], notes: list[str] +) -> None: + title = spec.get("metric_title") + if title: + created = {i: m for i, m in entities.list("metrics").items() if i not in before["metrics"]} + mine = {i: m for i, m in created.items() if (m.get("attributes") or {}).get("title") == title} + if not mine: + notes.append(f"OBSERVED no metric titled {title!r} was created") + for metric_id, metric in mine.items(): + maql = ((metric.get("attributes") or {}).get("content") or {}).get("maql") + notes.append(f"OBSERVED metric {metric_id} created, maql={maql!r}") + _delete(entities, "metrics", metric_id, notes) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_sinks.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_sinks.py new file mode 100644 index 000000000..0adc71215 --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_obfuscation_sinks.py @@ -0,0 +1,172 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +"""Read-back of the two sinks the agentic_obfuscation kind asserts on. + +* conversation DB -- the stored copy, read over the same API the UI uses on reload. +* Langfuse -- every trace of the conversation's session, with full observation input/output. + +Unlike the trace linking in ``_langfuse``, nothing here swallows an error: an absence +assertion over a sink that was never read would pass for the wrong reason, so every failure +to read raises :class:`SinkUnavailableError`. +""" + +from __future__ import annotations + +import os +import time +from collections.abc import Callable, Sequence +from typing import Any + +import httpx + +from gooddata_eval.core.agentic._obfuscation_check import contains + +_LANGFUSE_TIMEOUT_ENV = "GD_EVAL_OBFUSCATION_LANGFUSE_TIMEOUT_SEC" +_DEFAULT_LANGFUSE_TIMEOUT_SEC = 60.0 +TRACE_NOT_FOUND = "TRACE_NOT_FOUND" +_DB_TIMEOUT_SEC = 60.0 +_POLL_INTERVAL_SEC = 5.0 + + +class SinkUnavailableError(AssertionError): + """A sink could not be read, so nothing may be concluded from its contents.""" + + +class _LangfuseReadError(SinkUnavailableError): + """Langfuse answered a read with an error: the sink was not read, which is not an empty sink.""" + + +class TraceNotFoundError(SinkUnavailableError): + """Langfuse holds no trace for the conversation's session once the wait is over. + + Kept apart from a trace that arrived incomplete: it points at the export leg (a dropped + batch, a disabled trace flag, the wrong Langfuse project), not at masking. + """ + + +def _poll( + read: Callable[[], Any], + done: Callable[[Any], bool], + timeout_sec: float, + what: str, + shape: Callable[[Any], Any] = lambda doc: doc, + sleep: Callable[[float], None] = time.sleep, +) -> Any: + """Read until ``done`` holds and one further read has the same ``shape``. + + On timeout a non-empty document is still returned, so the verdict can say why it is + incomplete (a missing anchor, a redacted trace); only an empty one raises. + """ + deadline = time.monotonic() + timeout_sec + previous: Any = None + while True: + current = read() + if done(current) and previous is not None and shape(current) == previous: + return current + previous = shape(current) if done(current) else None + if time.monotonic() >= deadline: + if done(current) or current: + return current + raise SinkUnavailableError(f"{what}: not complete after {timeout_sec:.0f}s") + sleep(_POLL_INTERVAL_SEC) + + +def read_conversation_db( + http: httpx.Client, + conversations_url: str, + conversation_id: str, + anchors: Sequence[str], +) -> dict[str, Any]: + """The stored conversation (record + every item), once every turn's anchor is persisted. + + ``http`` carries the GoodData credentials. ``anchors`` may be empty for a rejected turn, + where nothing is expected to be stored. + """ + base = f"{conversations_url.rstrip('/')}/{conversation_id}" + + def get(url: str) -> Any: + try: + resp = http.get(url, headers={"Accept": "application/json"}) + except httpx.TimeoutException: + raise + except httpx.HTTPError as exc: + raise SinkUnavailableError(f"GET {url}: {exc}") from exc + if resp.status_code >= 400: + raise SinkUnavailableError(f"GET {url} -> {resp.status_code}: {resp.text[:300]}") + try: + return resp.json() + except ValueError as exc: + raise SinkUnavailableError(f"GET {url}: the body is not JSON") from exc + + def read() -> dict[str, Any] | None: + try: + items = get(f"{base}/items").get("items", []) + if not items and anchors: + return None + return {"conversation": get(base), "items": items} + except httpx.TimeoutException: + # A slow read is not a missing conversation: poll again until the deadline decides. + return None + + return _poll( + read, + lambda doc: doc is not None and all(contains(doc["items"], a) for a in anchors), + _DB_TIMEOUT_SEC, + f"conversation {conversation_id} read-back", + ) + + +def read_langfuse_session( + langfuse: Any, + conversation_id: str, + anchors: Sequence[str], + *, + allow_empty: bool = False, +) -> list[dict[str, Any]] | None: + """Every trace of the session with its observations, once every turn's anchor shows up. + + ``langfuse`` is an ``HttpxLangfuseClient``. ``allow_empty`` is for a rejected turn, which + may export no trace; it then returns ``None`` after the timeout. + """ + if langfuse is None: + raise SinkUnavailableError("no Langfuse client: set LANGFUSE_PUBLIC_KEY / LANGFUSE_SECRET_KEY") + timeout_sec = float(os.environ.get(_LANGFUSE_TIMEOUT_ENV) or _DEFAULT_LANGFUSE_TIMEOUT_SEC) + timed_out = 0 + + def read() -> list[dict[str, Any]]: + nonlocal timed_out + try: + return [langfuse.get_trace(trace_id) for trace_id in langfuse.session_trace_ids(conversation_id)] + except httpx.TimeoutException: + # A slow read is not a missing trace: poll again until the deadline decides. + timed_out += 1 + return [] + except httpx.HTTPError as exc: + raise _LangfuseReadError(f"Langfuse session {conversation_id}: {exc}") from exc + except (ValueError, LookupError, TypeError) as exc: + # A body that is not JSON, a row without an id, a listing past its page cap: read, but + # not readable. + raise _LangfuseReadError(f"Langfuse session {conversation_id}: unreadable response ({exc!r})") from exc + + def done(traces: list[dict[str, Any]]) -> bool: + return bool(traces) and all(contains(traces, a) for a in anchors) + + try: + # Latency and cost keep updating on a finished trace, so "settled" means no new trace + # and no new observation since the last read. + return _poll( + read, + done, + timeout_sec, + f"Langfuse session {conversation_id}", + shape=lambda traces: [(t.get("id"), len(t.get("observations") or [])) for t in traces], + ) + except _LangfuseReadError: + # A failed read says nothing about the export: never "no trace", never "none expected". + raise + except SinkUnavailableError as exc: + if allow_empty: + return None + slow = f" ({timed_out} read(s) timed out)" if timed_out else "" + raise TraceNotFoundError( + f"{TRACE_NOT_FOUND}: no Langfuse trace for session {conversation_id} after {timeout_sec:.0f}s{slow}" + ) from exc diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py index de77ca445..b4c023bc2 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/_trace_linker.py @@ -174,9 +174,10 @@ def observe(self, trace: Any, run_idx: int, *, conversation_id: str | None = Non output=output, ) - def score(self, trace_id: Any, *, name: str, value: Any, data_type: str) -> None: + def score(self, trace_id: Any, *, name: str, value: Any, data_type: str, comment: str | None = None) -> None: """Write one score, swallowing Langfuse failures the way ``score_safe`` always has.""" - self._lf.score_safe(self._client, trace_id, name=name, value=value, data_type=data_type) + extra = {"comment": comment} if comment else {} + self._lf.score_safe(self._client, trace_id, name=name, value=value, data_type=data_type, **extra) def quality(self, trace_id: Any, *, strict_checks: dict, latency_sec: Any, cost_usd: Any) -> None: """Write the derived quality/value scores for one run.""" diff --git a/packages/gooddata-eval/src/gooddata_eval/core/agentic/obfuscation.py b/packages/gooddata-eval/src/gooddata_eval/core/agentic/obfuscation.py new file mode 100644 index 000000000..027c96bbc --- /dev/null +++ b/packages/gooddata-eval/src/gooddata_eval/core/agentic/obfuscation.py @@ -0,0 +1,566 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +"""Agentic data-obfuscation evaluation runner. + +gen-ai masks sensitive values out of the copy it stores and the copy it exports to +Langfuse; what the model and the user's stream see stays as typed. So this kind does not +grade the answer for masking. It plants synthetic canaries in one or more user turns, reads +back the stored conversation and every Langfuse trace of the session, and decides by exact +substring (``_obfuscation_check``) -- no LLM decides whether a value leaked. + +``expected_output``: + canaries [{nonce, value, class, absent_from, present_in, mask_marker_present, + record_only_paths, ...}] -- see ``_obfuscation_check`` + status "enforced" | "known_limitation" + expected_turn_rejected {status_code, reason} when the first turn must be refused + observe optional record-only checks, see ``_obfuscation_observe`` + anchors optional; one non-sensitive fragment per turn, derived when absent + +A leak in any run fails the item whatever the gate: one stored secret is one too many. Every +other failure -- a missing trace, a blind sink, an over-masked value -- is a failed run the gate +decides on, as for every other kind. + +Before the first item of a workspace a preflight probe proves the environment can show a +result at all -- the database leg masks, the Langfuse leg masks, and both read-backs see +the probe. Without it an item could only pass by not looking. +""" + +from __future__ import annotations + +import json +import threading +import time +import uuid +from dataclasses import dataclass, field +from typing import Any + +import httpx + +from gooddata_eval.core._output import emit_line +from gooddata_eval.core.agentic._gate import ( + DEFAULT_GATE, + EvalGate, + gate_failure_note, + gate_passed, + log_gate_scores, + stamp_gate_metadata, +) +from gooddata_eval.core.agentic._obfuscation_check import ( + SINK_CONVERSATION_DB, + SINK_LANGFUSE_TRACE, + canary_needles, + derive_anchor, + evaluate_case, +) +from gooddata_eval.core.agentic._obfuscation_observe import WorkspaceEntities, observe_and_clean, snapshot +from gooddata_eval.core.agentic._obfuscation_sinks import ( + SinkUnavailableError, + read_conversation_db, + read_langfuse_session, +) +from gooddata_eval.core.agentic._trace_linker import ( + RunIdentity, + RunTraceContext, + SubmitTraceLink, + open_trace_window, + run_trace_link_inline, + submit_trace_scoring, + utc_now, +) +from gooddata_eval.core.chat.render import render_answer_text +from gooddata_eval.core.chat.sse_client import ChatClient, ChatError +from gooddata_eval.core.config import ReasoningEffort +from gooddata_eval.core.models import AgenticAssertionError, AgenticEvalOutcome, ChatResult + +_DEFAULT_K = 1 +_LOG = "[obfuscation]" +_PREFLIGHT_ATTEMPTS = 3 +_PREFLIGHT_RETRY_SEC = 10.0 +# User ids this short and common are scrubbed out of traces as plain words, which breaks up +# a canary such as "...@northwind-demo.invalid" before the detector sees it. +_SCRUB_PRONE_USER_IDS = frozenset({"demo", "test", "admin", "user"}) +_COMMENT_LIMIT = 2000 + +_preflight_lock = threading.Lock() +_preflight_results: dict[tuple[str, str], str | None] = {} + + +class ObfuscationAssertionError(AgenticAssertionError): + """Raised when an obfuscation evaluation fails: a leak, an over-masked value, or a blind sink.""" + + +@dataclass +class ObfuscationRun: + """One conversation of an obfuscation item: what was sent, what came back, the verdict.""" + + conversation_id: str + anchors: list[str] + # Per turn: None when answered, else {"status_code", "reason"} of the refusal. + errors: list[dict[str, Any] | None] = field(default_factory=list) + answers: list[str] = field(default_factory=list) + # Per turn: streamed text and visualizations, for record-only stream markers. + streams: list[str] = field(default_factory=list) + failures: list[str] = field(default_factory=list) + notes: list[str] = field(default_factory=list) + last_result: ChatResult | None = None + # Sinks whose read-back came back (a rejected turn's empty Langfuse leg counts): a verdict + # about a sink says nothing unless the sink is here. + sinks_read: set[str] = field(default_factory=set) + + @property + def trace_found(self) -> bool: + return SINK_LANGFUSE_TRACE in self.sinks_read + + @property + def leak_free(self) -> bool: + return bool(self.sinks_read) and not any(f.startswith("LEAK") for f in self.failures) + + @property + def passed(self) -> bool: + return not self.failures + + +@dataclass +class AgenticObfuscationSummary: + run_results: list[ObfuscationRun] + pass_at_k: bool + pass_power_k: bool + best: ObfuscationRun + + +def _anchors(expected: dict[str, Any], turns: list[str]) -> list[str]: + declared = expected.get("anchors") + if declared is not None: + return list(declared) + canaries = expected.get("canaries", []) + anchors = [] + for index, turn in enumerate(turns): + anchor = derive_anchor(turn, canaries) + if anchor is None: + raise ValueError( + f"turn {index}: no non-sensitive fragment of 12+ word characters to anchor the read-back on; " + "declare expected_output.anchors" + ) + anchors.append(anchor) + return anchors + + +def _send_turn(client: ChatClient, conversation_id: str, turn: str) -> tuple[ChatResult | None, dict[str, Any] | None]: + """One turn; a refused turn is returned, not raised, because some items expect it.""" + try: + return client.send_message(conversation_id, turn), None + except ChatError as exc: + return exc.partial_result, {"status_code": exc.status_code, "reason": exc.reason} + + +def _stream_of(result: ChatResult | None) -> str: + """What the turn streamed to the user: text, visualizations and alert proposals.""" + if result is None: + return "" + shown = result.model_dump(include={"text_response", "created_visualizations", "alert_proposals"}) + return json.dumps(shown, default=str) + + +def run_obfuscation_case( + client: ChatClient, + http: httpx.Client, + langfuse: Any, + conversations_url: str, + conversation_id: str, + turns: list[str], + expected: dict[str, Any], + entities: WorkspaceEntities | None = None, +) -> ObfuscationRun: + """Drive the turns, read back both sinks and judge them; report and clean up ``observe``.""" + spec = expected.get("observe") or {} + before = snapshot(entities, spec) if entities is not None and spec else None + run: ObfuscationRun | None = None + try: + run = _drive_and_judge(client, http, langfuse, conversations_url, conversation_id, turns, expected) + return run + finally: + # Also on failure: a persistent workspace must not keep what a failed run created. + if entities is not None and before is not None: + notes = observe_and_clean(entities, spec, before, run.streams[-1] if run and run.streams else "") + if run is not None: + run.notes.extend(notes) + + +def _drive_and_judge( + client: ChatClient, + http: httpx.Client, + langfuse: Any, + conversations_url: str, + conversation_id: str, + turns: list[str], + expected: dict[str, Any], +) -> ObfuscationRun: + rejection = expected.get("expected_turn_rejected") + run = ObfuscationRun(conversation_id, [] if rejection else _anchors(expected, turns)) + + for index, turn in enumerate(turns): + result, error = _send_turn(client, conversation_id, turn) + run.last_result = result or run.last_result + run.errors.append(error) + run.answers.append(render_answer_text(result) if result is not None else "") + run.streams.append(_stream_of(result)) + if error is not None: + break + if rejection: + run.failures.append(f"turn {index} was expected to be refused with {rejection} but was answered") + return run + + error = run.errors[-1] + if error is not None and not rejection: + hint = ( + " (503 DATA_OBFUSCATION_UNAVAILABLE: the inference gateway was not reachable)" + if error.get("reason") == "DATA_OBFUSCATION_UNAVAILABLE" + else "" + ) + run.failures.append( + f"turn {len(run.errors) - 1} failed with {error}, so the item says nothing about masking{hint}" + ) + return run + if rejection and ( + (error or {}).get("status_code") != rejection.get("status_code") + or (error or {}).get("reason") != rejection.get("reason") + ): + run.failures.append(f"turn refused with {error}, expected {rejection}") + + judge_sinks(run, http, langfuse, conversations_url, expected, rejected=bool(rejection)) + return run + + +def judge_sinks( + run: ObfuscationRun, + http: httpx.Client, + langfuse: Any, + conversations_url: str, + expected: dict[str, Any], + *, + rejected: bool = False, +) -> None: + """Read back both sinks of ``run``'s conversation and add their verdict to ``run``.""" + sinks: dict[str, Any] = {} + try: + sinks[SINK_CONVERSATION_DB] = read_conversation_db(http, conversations_url, run.conversation_id, run.anchors) + run.sinks_read.add(SINK_CONVERSATION_DB) + except SinkUnavailableError as exc: + sinks[SINK_CONVERSATION_DB] = None + run.failures.append(str(exc)) + try: + sinks[SINK_LANGFUSE_TRACE] = read_langfuse_session( + langfuse, run.conversation_id, run.anchors, allow_empty=rejected + ) + run.sinks_read.add(SINK_LANGFUSE_TRACE) + except SinkUnavailableError as exc: + sinks[SINK_LANGFUSE_TRACE] = None + run.failures.append(str(exc)) + + verdict = evaluate_case(expected, sinks, run.anchors, turn_rejected=rejected) + run.failures.extend(f for f in verdict.failures if f not in run.failures) + run.notes.extend(verdict.notes) + + +# The probe proves the environment masks at all, not that every stored copy is masked -- the +# items decide that. In the database only the stored message text is gated; a copy the agent +# puts elsewhere (a tool's step detail, conversation state) is reported, so a leak there +# fails the item that finds it instead of stopping every item of the run. Langfuse paths +# start with "$[" and stay gated everywhere. +_PROBE_RECORD_ONLY = [r"^\$\.(?!items\[\d+\]\.content\.text$)"] + + +def _probe() -> tuple[list[str], dict[str, Any]]: + nonce = f"preflight{uuid.uuid4().hex[:8]}" + # Worded so the agent answers without a tool: a value it searches for is copied into the + # step detail of that call. + question = ( + "Reply with the single word OK. Do not search, look anything up or use any tools. " + f"Reference only: qa.{nonce}@example.invalid" + ) + return [question], { + "canaries": [ + { + "nonce": nonce, + "value": f"qa.{nonce}@example.invalid", + "class": "EMAIL", + "absent_from": [SINK_CONVERSATION_DB, SINK_LANGFUSE_TRACE], + "mask_marker_present": "[EMAIL]", + "record_only_paths": _PROBE_RECORD_ONLY, + } + ], + "status": "enforced", + } + + +def _warn_on_scrub_prone_user(http: httpx.Client, host: str) -> None: + resp = http.get(f"{host.rstrip('/')}/api/v1/profile") + if resp.status_code >= 400: + emit_line(f"{_LOG} warning: /api/v1/profile answered {resp.status_code}; user-id scrub check skipped") + return + user_id = str(resp.json().get("userId", "")) + if user_id.lower() in _SCRUB_PRONE_USER_IDS: + emit_line( + f"{_LOG} warning: user id {user_id!r} is a common word; the Langfuse user-id scrub can split canaries " + "containing it and report leaks that are not the detector's" + ) + + +def _preflight(client: ChatClient, http: httpx.Client, langfuse: Any, host: str, conversations_url: str) -> str | None: + """None when the environment can show a masking result, else why it cannot.""" + _warn_on_scrub_prone_user(http, host) + last: ObfuscationRun | None = None + for attempt in range(1, _PREFLIGHT_ATTEMPTS + 1): + turns, expected = _probe() + conversation_id = client.create_conversation() + try: + last = run_obfuscation_case(client, http, langfuse, conversations_url, conversation_id, turns, expected) + finally: + client.delete_conversation(conversation_id) + if last.passed: + emit_line(f"{_LOG} preflight passed (attempt {attempt})") + for note in last.notes: + emit_line(f"{_LOG} preflight note: {note}") + return None + # The workspace setting reaches gen-ai through metadata sync, so a first probe can + # still run unmasked; a later one decides. + emit_line(f"{_LOG} preflight attempt {attempt} failed: {last.failures}") + if attempt < _PREFLIGHT_ATTEMPTS: + time.sleep(_PREFLIGHT_RETRY_SEC) + assert last is not None + return ( + f"obfuscation preflight failed: {last.failures}. Check: inference-gateway deployed; flags " + "enableGenAiDataObfuscation (database) and enableGenAiTraceObfuscation (Langfuse) on for the org; setting " + "enableAiDataObfuscation on for the workspace; LANGFUSE_* credentials of the project this environment " + "exports to." + ) + + +def require_preflight( + client: ChatClient, http: httpx.Client, langfuse: Any, host: str, workspace_id: str, conversations_url: str +) -> None: + key = (host.rstrip("/"), workspace_id) + with _preflight_lock: + if key not in _preflight_results: + _preflight_results[key] = _preflight(client, http, langfuse, host, conversations_url) + reason = _preflight_results[key] + if reason is not None: + raise ObfuscationAssertionError(reason) + + +def run_agentic_obfuscation( + host: str, + token: str, + workspace_id: str, + question: str, + expected_output: dict[str, Any], + k: int = _DEFAULT_K, + turns: list[str] | None = None, + initial_conversation_id: str | None = None, + reasoning_effort: ReasoningEffort | None = None, + agent_id: str | None = None, + langfuse: Any = None, +) -> AgenticObfuscationSummary: + """Run the obfuscation evaluation K times and return a summary. + + ``turns`` is the scripted conversation; a single-turn item passes only ``question``. + ``langfuse`` must be an ``HttpxLangfuseClient``: the Langfuse leg is read, not just scored. + """ + if k < 1: + raise ValueError(f"an obfuscation item needs at least one run, got k={k}") + all_turns = list(turns) if turns else [question] + client = ChatClient( + host=host, token=token, workspace_id=workspace_id, reasoning_effort=reasoning_effort, agent_id=agent_id + ) + http = httpx.Client(headers={"Authorization": f"Bearer {token}"}, timeout=60.0) + conversations_url = f"{host.rstrip('/')}/api/v1/ai/workspaces/{workspace_id}/chat/conversations" + entities = WorkspaceEntities(http, host, workspace_id) + run_results: list[ObfuscationRun] = [] + try: + require_preflight(client, http, langfuse, host, workspace_id, conversations_url) + for index in range(k): + owned = not (index == 0 and initial_conversation_id is not None) + conversation_id = client.create_conversation() if owned else initial_conversation_id + assert conversation_id is not None + try: + run = run_obfuscation_case( + client, http, langfuse, conversations_url, conversation_id, all_turns, expected_output, entities + ) + run_results.append(run) + emit_run(run, expected_output) + finally: + if owned: + client.delete_conversation(conversation_id) + finally: + client.close() + http.close() + + pass_at_k = any(r.passed for r in run_results) + pass_power_k = bool(run_results) and all(r.passed for r in run_results) + best = next((r for r in run_results if r.passed), run_results[0]) + return AgenticObfuscationSummary(run_results, pass_at_k, pass_power_k, best) + + +def emit_run(run: ObfuscationRun, expected: dict[str, Any]) -> None: + """Log one run: its turns, notes and verdict.""" + status = expected.get("status", "enforced") + emit_line(f"{_LOG} conversation={run.conversation_id} status={status} anchors={run.anchors}") + for index, (error, answer) in enumerate(zip(run.errors, run.answers)): + emit_line(f"{_LOG} turn {index}: " + (f"error={error}" if error else f"answer={answer[:200]!r}")) + for note in run.notes: + emit_line(f"{_LOG} note: {note}") + emit_line(f"{_LOG} verdict={'PASS' if run.passed else 'FAIL'}") + for failure in run.failures: + emit_line(f"{_LOG} - {failure}") + + +def without_canaries(text: str, canaries: list[dict[str, Any]]) -> str: + """Replace every canary spelling by its class, for text that leaves the evaluation (Langfuse).""" + for canary in canaries: + for needle in sorted(canary_needles(canary), key=len, reverse=True): + text = text.replace(needle, f"<{canary.get('class')}>") + return text + + +def evaluate_agentic_obfuscation( + host: str, + token: str, + workspace_id: str, + question: str, + expected_output: dict[str, Any], + k: int = _DEFAULT_K, + turns: list[str] | None = None, + initial_conversation_id: str | None = None, + agent_id: str | None = None, + langfuse: object | None = None, + dataset_item_id: str = "", + dataset_name: str = "obfuscation", + run_timestamp: str | None = None, + model_version_override: str | None = None, + run_metadata_extra: dict | None = None, + reasoning_effort: ReasoningEffort | None = None, + submit_trace_link: SubmitTraceLink = run_trace_link_inline, + gate: EvalGate = DEFAULT_GATE, +) -> AgenticEvalOutcome: + """Run the obfuscation evaluation, log to Langfuse, and raise ObfuscationAssertionError on failure.""" + if not isinstance(expected_output, dict) or not expected_output.get("canaries"): + raise ValueError("agentic_obfuscation expected_output must be an object with a non-empty 'canaries' list") + langfuse, window_start = open_trace_window(langfuse) + summary = run_agentic_obfuscation( + host=host, + token=token, + workspace_id=workspace_id, + question=question, + expected_output=expected_output, + k=k, + turns=turns, + initial_conversation_id=initial_conversation_id, + reasoning_effort=reasoning_effort, + agent_id=agent_id, + langfuse=langfuse, + ) + canaries = expected_output["canaries"] + leak_free = all(r.leak_free for r in summary.run_results) + leaked = any(f.startswith("LEAK") for r in summary.run_results for f in r.failures) + failures = [f"run {i}: {f}" for i, r in enumerate(summary.run_results) for f in r.failures] + item_passed = not leaked and gate_passed(gate, pass_at_k=summary.pass_at_k, pass_power_k=summary.pass_power_k) + + if langfuse is not None and dataset_item_id: + window_end = utc_now() + + def _write_scores(ctx: RunTraceContext) -> None: + stamp_gate_metadata(ctx.run_metadata, k=len(summary.run_results), gate=gate) + for run_idx, run in enumerate(summary.run_results): + pt = ctx.trace(run.conversation_id) + with ctx.observe( + pt, run_idx, conversation_id=run.conversation_id, output={"obfuscation_pass": run.passed} + ) as tid: + # The run's own failures: another run's leak must not annotate a run that passed. + comment = without_canaries("; ".join(run.failures), canaries)[:_COMMENT_LIMIT] or "no failure" + ctx.score( + tid, name="obfuscation_pass", value=float(run.passed), data_type="BOOLEAN", comment=comment + ) + ctx.score(tid, name="obfuscation_no_leak", value=float(run.leak_free), data_type="BOOLEAN") + ctx.score(tid, name="obfuscation_trace_found", value=float(run.trace_found), data_type="BOOLEAN") + recorded = [n for n in run.notes if n.startswith(("RECORDED", "OBSERVED"))] + if recorded: + # Record-only: what the SC does not rule on yet, next to the verdict it did not change. + ctx.score( + tid, + name="obfuscation_record_only_hit", + value=float(any(n.startswith("RECORDED") for n in recorded)), + data_type="BOOLEAN", + comment=without_canaries("; ".join(recorded), canaries)[:_COMMENT_LIMIT], + ) + # A leak in any run fails the item whatever the gate, so gate_passed -- the + # verdict reports read first -- must not say pass@K passed when one run leaked. + log_gate_scores( + ctx, + tid, + gate=gate, + pass_at_k=summary.pass_at_k and not leaked, + pass_power_k=summary.pass_power_k and not leaked, + ) + ctx.quality( + tid, + strict_checks={"obfuscation_no_leak": run.leak_free, "obfuscation_pass": run.passed}, + latency_sec=pt.latency if pt else None, + cost_usd=pt.total_cost if pt else None, + ) + + # Before the raise: a failing item's scores are the ones worth having. + submit_trace_scoring( + submit_trace_link, + RunIdentity( + host, + token, + workspace_id, + dataset_name, + run_timestamp, + model_version_override, + run_metadata_extra, + reasoning_effort, + ), + langfuse=langfuse, + dataset_item_id=dataset_item_id, + conversation_ids=[r.conversation_id for r in summary.run_results], + window_start=window_start, + window_end=window_end, + suffix_runs=len(summary.run_results) > 1, + write_scores=_write_scores, + item_input=turns or question, + ) + + best = summary.best + runs_passed = sum(1 for r in summary.run_results if r.passed) + runs_effective = len(summary.run_results) + last = best.last_result + detail = { + "leak_free": leak_free, + # False when the Langfuse leg was never read in some run: a no-leak verdict there would + # mean nothing. Named for the good outcome, as every boolean in a detail is read as a check. + "trace_found": all(r.trace_found for r in summary.run_results), + "failures": failures, + "notes": [n for r in summary.run_results for n in r.notes], + "actual_output": best.answers[-1] if best.answers else "", + } + reasoning_steps = list(last.reasoning_steps or []) if last else [] + response_id = last.response_id if last else None + + if not item_passed: + note = gate_failure_note(gate, runs_passed, runs_effective) + exc = ObfuscationAssertionError(f"Obfuscation assertion failed. {note} " + "; ".join(failures)) + exc.reasoning_steps = reasoning_steps + exc.conversation_id = best.conversation_id + exc.response_id = response_id + exc.detail = detail + exc.runs_passed = runs_passed + exc.runs_effective = runs_effective + raise exc + return AgenticEvalOutcome( + runs_passed=runs_passed, + runs_effective=runs_effective, + reasoning_steps=reasoning_steps, + conversation_id=best.conversation_id, + response_id=response_id, + detail=detail, + ) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py index 2127c88e3..00e72c29a 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/chat/sse_client.py @@ -72,11 +72,15 @@ def __init__( status_code: int | None = None, detail: str | None = None, partial_result: ChatResult | None = None, + reason: str | None = None, ) -> None: super().__init__(message) self.status_code = status_code self.detail = detail self.partial_result = partial_result + # gen-ai's machine-readable cause, e.g. DATA_OBFUSCATION_CONTENT_REJECTED; its data + # obfuscation errors carry this and no ``detail``. + self.reason = reason class TransientChatError(ChatError): @@ -382,12 +386,15 @@ def parse_sse_lines(lines: Iterable[str]) -> ChatResult: if "statusCode" in event_data: code = event_data.get("statusCode") detail = event_data.get("detail") - message = f"SSE error {code}: {detail}" + reason = event_data.get("reason") + message = f"SSE error {code}: {detail if detail is not None else reason}" if code in _RETRYABLE_STATUS_CODES: raise TransientChatError( - message, status_code=code, detail=detail, partial_result=_build_chat_result(acc) + message, status_code=code, detail=detail, partial_result=_build_chat_result(acc), reason=reason ) - raise ChatError(message, status_code=code, detail=detail, partial_result=_build_chat_result(acc)) + raise ChatError( + message, status_code=code, detail=detail, partial_result=_build_chat_result(acc), reason=reason + ) if event_data.get("responseId") and not acc.response_id: acc.response_id = event_data["responseId"] item = event_data.get("item") diff --git a/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py b/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py index 6c81ca179..60378faf0 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/dataset/langfuse_source.py @@ -30,12 +30,26 @@ def _make_client() -> httpx.Client: def _question_from_input(raw_input: Any) -> str: + return _question_and_turns(raw_input)[0] + + +def _question_and_turns(raw_input: Any) -> tuple[str, list[str] | None]: + """The item's question, plus its turns when the input is a scripted conversation. + + A list of strings is a multi-turn conversation whose first turn is the question. A + ``{"query": ...}`` object is accepted beside ``{"question": ...}``: Langfuse parses a + string input that happens to be valid JSON into an object, so a question that IS a JSON + document -- a pasted config -- can only survive the round trip wrapped. + """ if isinstance(raw_input, str): - return raw_input + return raw_input, None + if isinstance(raw_input, list) and raw_input and all(isinstance(turn, str) for turn in raw_input): + return raw_input[0], list(raw_input) if isinstance(raw_input, dict): - question = raw_input.get("question") - if isinstance(question, str): - return question + for key in ("question", "query"): + question = raw_input.get(key) + if isinstance(question, str): + return question, None raise ValueError(f"Unsupported Langfuse item input shape: {raw_input!r}") @@ -99,6 +113,9 @@ def _infer_test_kind(expected_output: object, default: str, metadata: object = N # {"expected_outputs": [...]} → experimental multi-candidate agentic vis if isinstance(eo.get("expected_outputs"), list): return "agentic_visualization" + # {"canaries": [...]} → data-obfuscation leak checks + if isinstance(eo.get("canaries"), list): + return "agentic_obfuscation" return default @@ -107,11 +124,13 @@ def _item_from_raw(raw: dict, *, dataset_name: str, test_kind: str) -> DatasetIt # REST API returns camelCase: expectedOutput, not expected_output expected_output = raw.get("expectedOutput") or raw.get("expected_output") resolved_kind = _infer_test_kind(expected_output, test_kind, raw.get("metadata")) + question, turns = _question_and_turns(raw.get("input")) return DatasetItem( id=str(raw["id"]), dataset_name=raw.get("datasetName") or dataset_name, test_kind=resolved_kind, - question=_question_from_input(raw.get("input")), + question=question, + turns=turns, expected_output=expected_output, summary_input=_summary_input_from_raw(raw, expected_output), user_context=_user_context_from_raw(raw), diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py index 1441457ec..f3195e7ba 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/client.py @@ -9,7 +9,7 @@ import threading import time import uuid -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from typing import Any import httpx @@ -25,6 +25,19 @@ _SCORES_PATH = "/api/public/scores" _OTLP_PATH = "/api/public/otel/v1/traces" +# Trace reads get their own timeout: a session-filtered listing and a full observation +# payload both run well past the default, which is sized for score writes. +_TRACE_READ_TIMEOUT = 60.0 +# Pages of observation rows read at most per trace read; a conversation holds a handful of +# traces of a few dozen observations each. +_OBSERVATION_PAGE_CAP = 20 +_OBSERVATION_PAGE_SIZE = 1000 +# The v2 observations endpoint cuts every metadata value at this length unless its key is +# named in ``expandMetadata``. +_METADATA_CUT = 200 +# How far back a trace read looks. The v2 endpoint wants a bounded window, and the +# conversations read back are minutes old. +_TRACE_READ_WINDOW = timedelta(days=1) _MAX_SCORE_ATTEMPTS = 3 _DEFAULT_RETRY_DELAY = 0.5 @@ -198,6 +211,69 @@ def list_traces( self._http, from_time=from_time, to_time=to_time, limit=limit, session_id=session_id ) + def _get_json(self, path: str, params: dict[str, Any] | None = None, *, timeout: float | None = None) -> Any: + """GET with the score path's retry on throttling and server errors. Raises on any other failure.""" + kwargs: dict[str, Any] = {"params": params} if timeout is None else {"params": params, "timeout": timeout} + resp = self._http.get(path, **kwargs) + for _retry in range(_MAX_SCORE_ATTEMPTS - 1): + if not _is_retryable(resp): + break + time.sleep(_retry_delay(resp)) + resp = self._http.get(path, **kwargs) + resp.raise_for_status() + return resp.json() + + def _observation_rows(self, params: dict[str, Any]) -> list[dict[str, Any]]: + """Every observation row matching ``params`` over a bounded window, all cursor pages.""" + now = datetime.now(timezone.utc) + base = { + "fromStartTime": (now - _TRACE_READ_WINDOW).isoformat(), + "toStartTime": (now + timedelta(minutes=5)).isoformat(), + "limit": _OBSERVATION_PAGE_SIZE, + **params, + } + rows: list[dict[str, Any]] = [] + cursor: str | None = None + for _page in range(_OBSERVATION_PAGE_CAP): + query = base if cursor is None else {**base, "cursor": cursor} + body = self._get_json(observations.OBSERVATIONS_PATH, query, timeout=_TRACE_READ_TIMEOUT) + rows.extend(body.get("data") or []) + cursor = (body.get("meta") or {}).get("cursor") + if not cursor: + return rows + raise LookupError(f"more than {_OBSERVATION_PAGE_CAP} pages of observations for {params}") + + def session_trace_ids(self, session_id: str) -> list[str]: + """Ids of every trace in one session, oldest first. gen-ai sets sessionId = conversationId.""" + first_seen: dict[str, str] = {} + for row in self._observation_rows({"sessionId": session_id, "fields": "core"}): + trace_id = row.get("traceId") + if trace_id: + start = row.get("startTime") or "" + first_seen[trace_id] = min(first_seen.get(trace_id, start), start) + return sorted(first_seen, key=lambda trace_id: first_seen[trace_id]) + + def get_trace(self, trace_id: str) -> dict[str, Any]: + """One trace as ``{"id", "observations"}``, every observation with its input, output and metadata. + + A metadata value the endpoint cut is read again in full: a canary past the cut would + otherwise pass unseen. + """ + params = {"traceId": trace_id, "fields": "core,basic,io,metadata"} + rows = self._observation_rows(params) + cut = sorted( + { + key + for row in rows + for key, value in (row.get("metadata") or {}).items() + if isinstance(value, str) and len(value) >= _METADATA_CUT + } + ) + if cut: + rows = self._observation_rows({**params, "expandMetadata": ",".join(cut)}) + rows.sort(key=lambda row: row.get("startTime") or "") + return {"id": trace_id, "observations": rows} + def flush(self) -> None: pass # no client-side batching diff --git a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py index 51686b8fe..b95d6e425 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/langfuse/observations.py @@ -9,7 +9,7 @@ if TYPE_CHECKING: import httpx -_OBSERVATIONS_PATH = "/api/public/v2/observations" +OBSERVATIONS_PATH = "/api/public/v2/observations" # Everything a TraceSummary needs: core (ids, times, parent), basic (sessionId), usage # (totalCost), metrics (latency), metadata. _FIELDS = "core,basic,usage,metrics,metadata" @@ -113,7 +113,7 @@ def list_traces_in_window( summaries: list[TraceSummary] = [] cursor: str | None = None for _page in range(max_pages): - resp = http.get(_OBSERVATIONS_PATH, params=params if cursor is None else {**params, "cursor": cursor}) + resp = http.get(OBSERVATIONS_PATH, params=params if cursor is None else {**params, "cursor": cursor}) resp.raise_for_status() body = resp.json() rows.extend(body.get("data") or []) diff --git a/packages/gooddata-eval/src/gooddata_eval/core/models.py b/packages/gooddata-eval/src/gooddata_eval/core/models.py index a3459fb01..c2de96773 100644 --- a/packages/gooddata-eval/src/gooddata_eval/core/models.py +++ b/packages/gooddata-eval/src/gooddata_eval/core/models.py @@ -8,7 +8,7 @@ import re from typing import Any -from pydantic import BaseModel, ConfigDict, Field, field_validator +from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator from gooddata_eval.core.timing import PhaseTimings @@ -427,3 +427,16 @@ class DatasetItem(BaseModel): # that schema (a discriminated union of view/widget descriptors), so re-modelling it here # would only create a second copy to keep in sync. user_context: dict[str, Any] | None = None + # A scripted multi-turn conversation, sent turn by turn; ``question`` is its first turn. + # Only used by the `agentic_obfuscation` test kind; ignored by all others. + turns: list[str] | None = None + + @model_validator(mode="before") + @classmethod + def _split_turns(cls, data: Any) -> Any: + """Accept a list of turns as ``question``, the shape a multi-turn fixture file carries.""" + if isinstance(data, dict): + question = data.get("question") + if isinstance(question, list) and question and all(isinstance(turn, str) for turn in question): + return {**data, "question": question[0], "turns": data.get("turns") or list(question)} + return data diff --git a/packages/gooddata-eval/tests/test_agentic_obfuscation.py b/packages/gooddata-eval/tests/test_agentic_obfuscation.py new file mode 100644 index 000000000..e55d96b19 --- /dev/null +++ b/packages/gooddata-eval/tests/test_agentic_obfuscation.py @@ -0,0 +1,581 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +# SPDX-License-Identifier: LicenseRef-GoodData-Enterprise +import json +import re +from contextlib import contextmanager +from typing import Any +from unittest.mock import MagicMock, patch + +import httpx +import pytest +from gooddata_eval.core.agentic import _obfuscation_sinks, obfuscation +from gooddata_eval.core.agentic._obfuscation_check import SINK_CONVERSATION_DB, SINK_LANGFUSE_TRACE +from gooddata_eval.core.agentic._obfuscation_observe import WorkspaceEntities, observe_and_clean, snapshot +from gooddata_eval.core.agentic._obfuscation_sinks import ( + SinkUnavailableError, + TraceNotFoundError, + read_conversation_db, + read_langfuse_session, +) +from gooddata_eval.core.agentic.obfuscation import ObfuscationAssertionError, evaluate_agentic_obfuscation +from gooddata_eval.core.chat.sse_client import ChatError +from gooddata_eval.core.models import ChatResult + +_EMAIL = "marie.canary7f3a@example.invalid" +_QUESTION = f"Can you check the recent orders for the account registered to {_EMAIL}?" +_MASKED = "Can you check the recent orders for the account registered to [EMAIL]?" +_EXPECTED: dict[str, Any] = { + "canaries": [ + { + "nonce": "canary7f3a", + "value": _EMAIL, + "class": "EMAIL", + "absent_from": [SINK_LANGFUSE_TRACE, SINK_CONVERSATION_DB], + "mask_marker_present": "[EMAIL]", + } + ], + "status": "enforced", +} + + +@pytest.fixture(autouse=True) +def _fresh_preflight(monkeypatch): + """Every test decides its own preflight outcome; none inherits a cached one.""" + monkeypatch.setattr(obfuscation, "_preflight_results", {}) + monkeypatch.setattr(obfuscation, "_PREFLIGHT_RETRY_SEC", 0) + monkeypatch.setattr(_obfuscation_sinks, "_POLL_INTERVAL_SEC", 0) + + +def _answer(text: str = "Here are the recent orders.") -> ChatResult: + return ChatResult.model_validate({"textResponse": text, "toolCallEvents": [], "responseId": "resp-1"}) + + +def _db(*texts: str) -> dict: + items = [{"role": "user", "content": {"type": "text", "text": t}} for t in texts] + return {"conversation": {"title": None}, "items": items} + + +def _trace(*texts: str) -> list: + return [ + {"id": f"t{i}", "input": json.dumps({"parts": [{"text": t}]}), "observations": []} for i, t in enumerate(texts) + ] + + +@contextmanager +def _patched(client, *, db, lf): + with ( + patch("gooddata_eval.core.agentic.obfuscation.ChatClient", return_value=client), + patch("gooddata_eval.core.agentic.obfuscation.read_conversation_db", side_effect=db), + patch("gooddata_eval.core.agentic.obfuscation.read_langfuse_session", side_effect=lf), + patch("gooddata_eval.core.agentic.obfuscation._warn_on_scrub_prone_user"), + ): + yield + + +class _Chat: + """A chat client whose conversations remember the turns sent to them.""" + + def __init__(self, fail_with: dict[str, ChatError] | None = None) -> None: + self.sent: dict[str, list[str]] = {} + self._n = 0 + self._fail_with = fail_with or {} + + def create_conversation(self) -> str: + self._n += 1 + conversation_id = f"conv-{self._n}" + self.sent[conversation_id] = [] + return conversation_id + + def send_message(self, conversation_id: str, question: str) -> ChatResult: + self.sent.setdefault(conversation_id, []).append(question) + for marker, exc in self._fail_with.items(): + if marker in question: + raise exc + return _answer() + + def delete_conversation(self, conversation_id: str) -> None: + pass + + def close(self) -> None: + pass + + +_ANY_EMAIL = re.compile(r"[\w.+-]+@[\w-]+(?:\.[\w-]+)+") + + +def _is_probe(chat: _Chat, conversation_id: str) -> bool: + return "@example.invalid" in chat.sent[conversation_id][0] and "qa.preflight" in chat.sent[conversation_id][0] + + +def _masking(chat: _Chat, *, leak_into_langfuse: bool = False): + """Sinks that store the masked copy; optionally let an item's (not the probe's) canary reach Langfuse.""" + + def masked(conversation_id: str) -> list[str]: + return [_ANY_EMAIL.sub("[EMAIL]", t) for t in chat.sent[conversation_id]] + + def db(_http, _url, conversation_id, anchors, **_kw): + return _db(*masked(conversation_id)) + + def lf(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + raw = leak_into_langfuse and not _is_probe(chat, conversation_id) + return _trace(*(chat.sent[conversation_id] if raw else masked(conversation_id))) + + return db, lf + + +def _evaluate(**kwargs: Any): + kwargs.setdefault("expected_output", _EXPECTED) + return evaluate_agentic_obfuscation(host="https://h", token="tok", workspace_id="ws1", question=_QUESTION, **kwargs) + + +def test_masked_in_both_sinks_passes(): + chat = _Chat() + db, lf = _masking(chat) + with _patched(chat, db=db, lf=lf): + outcome = _evaluate() + assert outcome.detail["leak_free"] is True + assert outcome.detail["trace_found"] is True + assert outcome.runs_passed == 1 + + +def test_canary_reaching_langfuse_fails_as_a_leak(): + chat = _Chat() + db, lf = _masking(chat, leak_into_langfuse=True) + with _patched(chat, db=db, lf=lf), pytest.raises(ObfuscationAssertionError, match="LEAK EMAIL"): + _evaluate() + + +def test_preflight_failure_stops_every_item_of_the_workspace_without_retrying_it(): + chat = _Chat() + + def unmasked_db(_http, _url, conversation_id, anchors, **_kw): + return _db(*chat.sent[conversation_id]) + + _, lf = _masking(chat) + with _patched(chat, db=unmasked_db, lf=lf): + with pytest.raises(ObfuscationAssertionError, match="preflight failed"): + _evaluate() + probes = len(chat.sent) + with pytest.raises(ObfuscationAssertionError, match="preflight failed"): + _evaluate() + assert len(chat.sent) == probes, "the cached verdict must not probe again" + assert all(_is_probe(chat, c) for c in chat.sent), "no item ran behind a failed preflight" + + +def test_a_probe_copied_into_a_tool_step_detail_does_not_fail_the_preflight(): + chat = _Chat() + db, lf = _masking(chat) + + def db_with_step_detail(_http, _url, conversation_id, anchors, **_kw): + stored = db(_http, _url, conversation_id, anchors) + if _is_probe(chat, conversation_id): + found = _ANY_EMAIL.search(chat.sent[conversation_id][0]) + assert found is not None + raw = found.group(0) + stored["items"].append({"role": "tool", "detail": {"category": "catalogSearch", "query": [raw]}}) + return stored + + with _patched(chat, db=db_with_step_detail, lf=lf): + outcome = _evaluate() + assert outcome.runs_passed == 1, "the item ran and passed behind the preflight" + + +def test_a_probe_reaching_langfuse_fails_the_preflight(): + chat = _Chat() + db, _ = _masking(chat) + + def raw_trace(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + return _trace(*chat.sent[conversation_id]) + + with _patched(chat, db=db, lf=raw_trace), pytest.raises(ObfuscationAssertionError, match="preflight failed"): + _evaluate() + + +def test_refused_turn_is_the_expected_outcome_when_the_item_says_so(): + rejected = ChatError("SSE error 422", status_code=422, reason="DATA_OBFUSCATION_CONTENT_REJECTED") + chat = _Chat(fail_with={"C:\\Users": rejected}) + expected = { + "canaries": [{**_EXPECTED["canaries"][0], "mask_marker_present": None}], + "expected_turn_rejected": {"status_code": 422, "reason": "DATA_OBFUSCATION_CONTENT_REJECTED"}, + } + db, lf = _masking(chat) + + def no_trace(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + return lf(_langfuse, conversation_id, anchors) if _is_probe(chat, conversation_id) else None + + with _patched(chat, db=db, lf=no_trace): + outcome = evaluate_agentic_obfuscation( + host="https://h", + token="tok", + workspace_id="ws1", + question=f"The log line was C:\\Users\\svc\\{_EMAIL} and it aborted", + expected_output=expected, + ) + assert outcome.detail["leak_free"] is True + + +def test_an_unexpected_refusal_says_nothing_about_masking(): + unavailable = ChatError("SSE error 503", status_code=503, reason="DATA_OBFUSCATION_UNAVAILABLE") + chat = _Chat(fail_with={"recent orders": unavailable}) + db, lf = _masking(chat) + with _patched(chat, db=db, lf=lf), pytest.raises(ObfuscationAssertionError, match="DATA_OBFUSCATION_UNAVAILABLE"): + _evaluate() + + +def test_turns_are_sent_in_order_to_one_conversation(): + chat = _Chat() + db, lf = _masking(chat) + turns = [_QUESTION, "Now show the same Total Sales by quarter instead."] + with _patched(chat, db=db, lf=lf): + _evaluate(turns=turns) + item_conversations = [t for c, t in chat.sent.items() if not _is_probe(chat, c)] + assert item_conversations == [turns] + + +def test_an_answer_rubric_does_not_decide_the_verdict(): + chat = _Chat() + db, lf = _masking(chat) + with _patched(chat, db=db, lf=lf): + outcome = _evaluate(expected_output={**_EXPECTED, "answer_rubric": "answer the analytics question"}) + assert outcome.runs_passed == 1 + + +def test_a_leak_in_one_run_fails_the_item_whatever_the_gate(): + chat = _Chat() + db, lf = _masking(chat) + calls = {"n": 0} + + def one_bad_run(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + if _is_probe(chat, conversation_id): + return lf(_langfuse, conversation_id, anchors) + calls["n"] += 1 + return _trace(*chat.sent[conversation_id]) if calls["n"] == 2 else lf(_langfuse, conversation_id, anchors) + + with _patched(chat, db=db, lf=one_bad_run), pytest.raises(ObfuscationAssertionError, match="run 1: LEAK"): + _evaluate(k=2, gate="any") + + +def test_a_leak_in_one_run_is_not_scored_as_a_passed_gate(): + chat = _Chat() + db, lf = _masking(chat) + calls = {"n": 0} + + def one_bad_run(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + if _is_probe(chat, conversation_id): + return lf(_langfuse, conversation_id, anchors) + calls["n"] += 1 + return _trace(*chat.sent[conversation_id]) if calls["n"] == 2 else lf(_langfuse, conversation_id, anchors) + + with ( + _patched(chat, db=db, lf=one_bad_run), + patch("gooddata_eval.core.agentic.obfuscation.submit_trace_scoring") as submit, + pytest.raises(ObfuscationAssertionError), + ): + _evaluate(k=2, gate="any", langfuse=MagicMock(), dataset_item_id="i1") + ctx = MagicMock() + ctx.observe.return_value.__exit__.return_value = False # let an error in the block surface + submit.call_args.kwargs["write_scores"](ctx) + gate = [c.kwargs["value"] for c in ctx.score.call_args_list if c.kwargs["name"] == "gate_passed"] + assert gate == [False, False], "reports read gate_passed first; it must agree with the failed item" + + +def _one_run_without_a_trace(chat, lf): + calls = {"n": 0} + + def read(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + if _is_probe(chat, conversation_id): + return lf(_langfuse, conversation_id, anchors) + calls["n"] += 1 + if calls["n"] == 1: + raise TraceNotFoundError(f"TRACE_NOT_FOUND: no Langfuse trace for session {conversation_id} after 60s") + return lf(_langfuse, conversation_id, anchors) + + return read + + +def test_a_run_without_a_trace_is_a_failed_run_the_any_gate_can_carry(): + chat = _Chat() + db, lf = _masking(chat) + with _patched(chat, db=db, lf=_one_run_without_a_trace(chat, lf)): + outcome = _evaluate(k=2, gate="any") + assert outcome.runs_passed == 1 + assert outcome.detail["trace_found"] is False, "the run that saw no trace is still reported" + + +def test_a_run_without_a_trace_fails_the_power_gate(): + chat = _Chat() + db, lf = _masking(chat) + with ( + _patched(chat, db=db, lf=_one_run_without_a_trace(chat, lf)), + pytest.raises(ObfuscationAssertionError, match="run 0: TRACE_NOT_FOUND"), + ): + _evaluate(k=2, gate="power") + + +def test_a_turn_that_failed_reads_no_sink_so_it_claims_no_clean_result(): + unavailable = ChatError("SSE error 503", status_code=503, reason="DATA_OBFUSCATION_UNAVAILABLE") + chat = _Chat(fail_with={"recent orders": unavailable}) + db, lf = _masking(chat) + with _patched(chat, db=db, lf=lf), pytest.raises(ObfuscationAssertionError) as raised: + _evaluate() + assert raised.value.detail["leak_free"] is False + assert raised.value.detail["trace_found"] is False + + +def test_each_run_is_scored_with_its_own_failures(): + chat = _Chat() + db, lf = _masking(chat) + calls = {"n": 0} + + def second_run_leaks(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + if _is_probe(chat, conversation_id): + return lf(_langfuse, conversation_id, anchors) + calls["n"] += 1 + return _trace(*chat.sent[conversation_id]) if calls["n"] == 2 else lf(_langfuse, conversation_id, anchors) + + with ( + _patched(chat, db=db, lf=second_run_leaks), + patch("gooddata_eval.core.agentic.obfuscation.submit_trace_scoring") as submit, + pytest.raises(ObfuscationAssertionError), + ): + _evaluate(k=2, gate="any", langfuse=MagicMock(), dataset_item_id="i1") + ctx = MagicMock() + ctx.observe.return_value.__exit__.return_value = False + submit.call_args.kwargs["write_scores"](ctx) + comments = [c.kwargs.get("comment") for c in ctx.score.call_args_list if c.kwargs["name"] == "obfuscation_pass"] + assert comments[0] == "no failure" + assert "LEAK EMAIL" in comments[1] and "canary7f3a" not in comments[1] + + +def test_expected_output_without_canaries_is_rejected(): + with pytest.raises(ValueError, match="canaries"): + evaluate_agentic_obfuscation( + host="https://h", token="tok", workspace_id="ws1", question="q", expected_output={"status": "enforced"} + ) + + +# --- sinks -------------------------------------------------------------------------- + + +def _gd_http(handler) -> httpx.Client: + return httpx.Client(transport=httpx.MockTransport(handler)) + + +def test_db_read_back_waits_for_the_anchor_and_returns_record_and_items(): + reads = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + if request.url.path.endswith("/items"): + reads["n"] += 1 + items = [] if reads["n"] == 1 else [{"role": "user", "content": {"text": _MASKED}}] + return httpx.Response(200, json={"items": items}) + return httpx.Response(200, json={"title": "t"}) + + with _gd_http(handler) as http: + doc = read_conversation_db(http, "https://h/api/v1/ai/workspaces/ws/chat/conversations", "c1", ["check the"]) + assert doc["conversation"] == {"title": "t"} + assert doc["items"][0]["content"]["text"] == _MASKED + + +def test_db_read_back_error_is_a_sink_error_not_an_empty_sink(): + with ( + _gd_http(lambda request: httpx.Response(403, text="forbidden")) as http, + pytest.raises(SinkUnavailableError, match="403"), + ): + read_conversation_db(http, "https://h/api/v1/ai/workspaces/ws/chat/conversations", "c1", ["x"]) + + +def test_langfuse_read_back_needs_a_client(): + with pytest.raises(SinkUnavailableError, match="no Langfuse client"): + read_langfuse_session(None, "c1", ["x"]) + + +def test_langfuse_read_back_returns_every_trace_of_the_session(): + langfuse = MagicMock() + langfuse.session_trace_ids.return_value = ["t1"] + langfuse.get_trace.return_value = {"id": "t1", "input": _MASKED, "observations": [{"id": "o1"}]} + traces = read_langfuse_session(langfuse, "c1", ["check the"]) + assert traces == [langfuse.get_trace.return_value] + langfuse.session_trace_ids.assert_called_with("c1") + + +def test_langfuse_read_back_of_a_rejected_turn_may_be_empty(monkeypatch): + monkeypatch.setattr(_obfuscation_sinks, "_DEFAULT_LANGFUSE_TIMEOUT_SEC", 0) + langfuse = MagicMock() + langfuse.session_trace_ids.return_value = [] + assert read_langfuse_session(langfuse, "c1", [], allow_empty=True) is None + with pytest.raises(TraceNotFoundError, match="^TRACE_NOT_FOUND: no Langfuse trace for session c1 after 0s$"): + read_langfuse_session(langfuse, "c1", ["x"]) + + +def test_a_langfuse_read_error_is_never_taken_for_a_turn_that_exported_nothing(monkeypatch): + monkeypatch.setattr(_obfuscation_sinks, "_DEFAULT_LANGFUSE_TIMEOUT_SEC", 0) + langfuse = MagicMock() + denied = httpx.Response(401, request=httpx.Request("GET", "https://lf/api/public/sessions/c1")) + langfuse.session_trace_ids.side_effect = httpx.HTTPStatusError("401", request=denied.request, response=denied) + with pytest.raises(SinkUnavailableError, match="Langfuse session c1") as raised: + read_langfuse_session(langfuse, "c1", [], allow_empty=True) + assert not isinstance(raised.value, TraceNotFoundError) + with pytest.raises(SinkUnavailableError) as raised: + read_langfuse_session(langfuse, "c1", ["x"]) + assert not isinstance(raised.value, TraceNotFoundError), "a failed read is not a missing trace" + + +def test_a_db_transport_or_body_error_is_a_sink_error(monkeypatch): + monkeypatch.setattr(_obfuscation_sinks, "_DB_TIMEOUT_SEC", 0) + url = "https://h/api/v1/ai/workspaces/ws/chat/conversations" + + def refused(request: httpx.Request) -> httpx.Response: + raise httpx.ConnectError("refused", request=request) + + with _gd_http(refused) as http, pytest.raises(SinkUnavailableError, match="refused"): + read_conversation_db(http, url, "c1", ["x"]) + with ( + _gd_http(lambda request: httpx.Response(200, text="")) as http, + pytest.raises(SinkUnavailableError, match="not JSON"), + ): + read_conversation_db(http, url, "c1", ["x"]) + + +def test_a_slow_db_read_is_polled_again(): + reads = {"n": 0} + item = {"role": "user", "content": {"type": "text", "text": "Can you check the recent orders"}} + + def handler(request: httpx.Request) -> httpx.Response: + reads["n"] += 1 + if reads["n"] == 1: + raise httpx.ReadTimeout("slow", request=request) + return httpx.Response(200, json={"items": [item]} if request.url.path.endswith("/items") else {"id": "c1"}) + + with _gd_http(handler) as http: + doc = read_conversation_db(http, "https://h/api/v1/ai/workspaces/ws/chat/conversations", "c1", ["check the"]) + assert doc["items"] == [item] + + +def test_an_item_needs_at_least_one_run(): + with pytest.raises(ValueError, match="at least one run"): + _evaluate(k=0) + + +@pytest.mark.parametrize("error", [ValueError("not JSON"), KeyError("id")]) +def test_an_unreadable_langfuse_response_is_a_sink_error(error): + langfuse = MagicMock() + langfuse.session_trace_ids.side_effect = error + with pytest.raises(SinkUnavailableError, match="unreadable response") as raised: + read_langfuse_session(langfuse, "c1", ["x"]) + assert not isinstance(raised.value, TraceNotFoundError) + + +def test_trace_not_found_says_how_many_reads_timed_out(monkeypatch): + monkeypatch.setattr(_obfuscation_sinks, "_DEFAULT_LANGFUSE_TIMEOUT_SEC", 0) + langfuse = MagicMock() + langfuse.session_trace_ids.side_effect = httpx.ReadTimeout("slow") + with pytest.raises(TraceNotFoundError, match=r"\(1 read\(s\) timed out\)$"): + read_langfuse_session(langfuse, "c1", ["x"]) + + +def test_the_default_wait_for_a_trace_is_sixty_seconds(): + assert _obfuscation_sinks._DEFAULT_LANGFUSE_TIMEOUT_SEC == 60.0 + + +def test_an_item_whose_trace_never_arrives_fails_marked_as_trace_not_found(): + chat = _Chat() + db, lf = _masking(chat) + + def no_item_trace(_langfuse, conversation_id, anchors, allow_empty=False, **_kw): + if _is_probe(chat, conversation_id): + return lf(_langfuse, conversation_id, anchors) + raise TraceNotFoundError(f"TRACE_NOT_FOUND: no Langfuse trace for session {conversation_id} after 60s") + + with _patched(chat, db=db, lf=no_item_trace), pytest.raises(ObfuscationAssertionError) as raised: + _evaluate() + assert "run 0: TRACE_NOT_FOUND" in str(raised.value) + assert raised.value.detail["trace_found"] is False + assert raised.value.detail["leak_free"] is True, "a missing trace is not a leak" + + +# --- observe ------------------------------------------------------------------------ + + +def test_observe_reports_and_deletes_only_the_automation_the_item_made(): + listed = {"n": 0} + deleted: list[str] = [] + pre = {"id": "old", "attributes": {"metadata": "total_sales"}} + made = { + "id": "new", + "attributes": {"metadata": "total_sales", "externalRecipients": [{"email": "[EMAIL]"}]}, + "relationships": {"recipients": {"data": [{"id": "u1"}]}}, + } + other = {"id": "someone-else", "attributes": {"metadata": "net_sales"}} + + def handler(request: httpx.Request) -> httpx.Response: + if request.method == "DELETE": + deleted.append(request.url.path.rsplit("/", 1)[1]) + return httpx.Response(204) + listed["n"] += 1 + return httpx.Response(200, json={"data": [pre] if listed["n"] == 1 else [pre, made, other]}) + + with _gd_http(handler) as http: + entities = WorkspaceEntities(http, "https://h", "ws") + spec = {"automation_match": "total_sales", "stream_markers": ["[EMAIL]"]} + before = snapshot(entities, spec) + notes = observe_and_clean(entities, spec, before, last_stream='{"text": "sent to [EMAIL]"}') + assert deleted == ["new"] + assert "OBSERVED automation new created with external=['[EMAIL]'] internal=['u1']" in notes + assert "OBSERVED last turn stream contains '[EMAIL]'" in notes + + +def test_observe_says_so_when_the_chat_created_nothing(): + with _gd_http(lambda request: httpx.Response(200, json={"data": []})) as http: + entities = WorkspaceEntities(http, "https://h", "ws") + spec = {"metric_title": "Average Order Value QA"} + notes = observe_and_clean(entities, spec, snapshot(entities, spec), last_stream="") + assert notes == ["OBSERVED no metric titled 'Average Order Value QA' was created"] + + +def test_observe_cleanup_errors_become_notes_and_do_not_stop_the_rest(): + made = [ + {"id": "a1", "attributes": {"metadata": "total_sales"}}, + {"id": "a2", "attributes": {"metadata": "total_sales"}}, + ] + listed = {"n": 0} + deleted: list[str] = [] + + def handler(request: httpx.Request) -> httpx.Response: + if request.method == "DELETE": + entity_id = request.url.path.rsplit("/", 1)[1] + deleted.append(entity_id) + return httpx.Response(404 if entity_id == "a1" else 500) + if "/metrics" in request.url.path: + return httpx.Response(503) + listed["n"] += 1 + return httpx.Response(200, json={"data": [] if listed["n"] == 1 else made}) + + with _gd_http(handler) as http: + entities = WorkspaceEntities(http, "https://h", "ws") + spec = {"automation_match": "total_sales"} + before = snapshot(entities, spec) + notes = observe_and_clean(entities, {**spec, "metric_title": "M"}, {**before, "metrics": set()}, "") + assert deleted == ["a1", "a2"], "a 404 is already gone and a 500 does not stop the next delete" + assert any(n.startswith("OBSERVED could not delete automations a2") for n in notes) + assert any(n.startswith("OBSERVED cleanup failed (metrics)") for n in notes) + + +def test_observe_lists_every_page_of_a_busy_workspace(monkeypatch): + monkeypatch.setattr("gooddata_eval.core.agentic._obfuscation_observe._PAGE_SIZE", 2) + pages = [[{"id": "a"}, {"id": "b"}], [{"id": "c"}, {"id": "d"}], [{"id": "e"}]] + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"data": pages[int(request.url.params["page"])]}) + + with _gd_http(handler) as http: + listed = WorkspaceEntities(http, "https://h", "ws").list("automations") + assert sorted(listed) == ["a", "b", "c", "d", "e"] + + +def test_a_slow_langfuse_read_is_polled_again_not_reported_as_missing(): + langfuse = MagicMock() + langfuse.session_trace_ids.return_value = ["t1"] + trace = {"id": "t1", "input": _MASKED, "observations": [{"id": "o1"}]} + langfuse.get_trace.side_effect = [httpx.ReadTimeout("slow"), trace, trace] + assert read_langfuse_session(langfuse, "c1", ["check the"]) == [trace] diff --git a/packages/gooddata-eval/tests/test_agentic_runner.py b/packages/gooddata-eval/tests/test_agentic_runner.py index 28ce3a4c1..0b36693be 100644 --- a/packages/gooddata-eval/tests/test_agentic_runner.py +++ b/packages/gooddata-eval/tests/test_agentic_runner.py @@ -91,6 +91,7 @@ def test_dispatch_agentic_omits_agent_id_by_default(): ("agentic_kda_skill", {"Measure": {"type": "metric", "id": "revenue"}}, "evaluate_agentic_kda_skill"), ("agentic_conversation", {"fixture": _MIN_CONVERSATION_FIXTURE}, "evaluate_agentic_conversation"), ("agentic_what_if", {"metric_id": "spend"}, "evaluate_agentic_what_if"), + ("agentic_obfuscation", {"canaries": [{"nonce": "n1", "value": "a@b.invalid"}]}, "evaluate_agentic_obfuscation"), ] diff --git a/packages/gooddata-eval/tests/test_cli.py b/packages/gooddata-eval/tests/test_cli.py index 1d23293bf..79c2253a0 100644 --- a/packages/gooddata-eval/tests/test_cli.py +++ b/packages/gooddata-eval/tests/test_cli.py @@ -28,6 +28,12 @@ def test_build_run_config_requires_a_source(): cli_main.parse_args(["run", "--host", "h", "--workspace", "w"]) +@pytest.mark.parametrize("runs", ["0", "-1", "two"]) +def test_parse_args_rejects_a_run_count_below_one(runs): + with pytest.raises(SystemExit): + cli_main.parse_args(["run", "--host", "h", "--workspace", "w", "--dataset", "d", "--runs", runs]) + + def test_parse_args_agent_id_flag(): args = cli_main.parse_args(["run", "--host", "h", "--workspace", "w", "--dataset", "d", "--agent-id", "agent-1"]) assert args.agent_id == "agent-1" diff --git a/packages/gooddata-eval/tests/test_langfuse_client.py b/packages/gooddata-eval/tests/test_langfuse_client.py index 931d7588e..62d717782 100644 --- a/packages/gooddata-eval/tests/test_langfuse_client.py +++ b/packages/gooddata-eval/tests/test_langfuse_client.py @@ -377,3 +377,96 @@ def test_a_dataset_run_item_for_an_unknown_item_raises(make_client): client = make_client(lambda request: httpx.Response(404, json={})) with pytest.raises(LookupError): client.api.dataset_run_items.create(run_name="run", dataset_item_id="local", trace_id="t-1") + + +def test_session_trace_ids_come_from_v2_observations_by_session_oldest_first(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + rows = [ + {"id": "o3", "traceId": "t2", "startTime": "2026-09-23T10:00:02Z"}, + {"id": "o1", "traceId": "t1", "startTime": "2026-09-23T10:00:01Z"}, + {"id": "o2", "traceId": "t1", "startTime": "2026-09-23T10:00:03Z"}, + ] + return httpx.Response(200, json={"data": rows, "meta": {"cursor": None}}) + + assert make_client(handler).session_trace_ids("conv-1") == ["t1", "t2"] + params = seen[0].url.params + assert seen[0].url.path == "/api/public/v2/observations" + assert params["sessionId"] == "conv-1" + assert params["fromStartTime"] and params["toStartTime"], "the v2 endpoint is read over a bounded window" + + +def test_observation_reads_follow_the_cursor(make_client): + pages = { + "": ([{"id": "o1", "traceId": "t1", "startTime": "1"}], "c2"), + "c2": ([{"id": "o2", "traceId": "t2", "startTime": "2"}], None), + } + + def handler(request: httpx.Request) -> httpx.Response: + rows, cursor = pages[request.url.params.get("cursor", "")] + return httpx.Response(200, json={"data": rows, "meta": {"cursor": cursor}}) + + assert make_client(handler).session_trace_ids("conv-1") == ["t1", "t2"] + + +def test_get_trace_returns_the_trace_observations_with_their_io(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + rows = [ + {"id": "o2", "traceId": "t1", "startTime": "2", "input": "[EMAIL]", "metadata": {}}, + {"id": "o1", "traceId": "t1", "startTime": "1", "input": "Can you check", "metadata": {}}, + ] + return httpx.Response(200, json={"data": rows, "meta": {}}) + + trace = make_client(handler).get_trace("t1") + assert trace["id"] == "t1" + assert [o["id"] for o in trace["observations"]] == ["o1", "o2"] + assert seen[0].url.params["traceId"] == "t1" + assert "io" in seen[0].url.params["fields"].split(",") + + +def test_get_trace_reads_a_cut_metadata_value_again_in_full(make_client): + seen: list[httpx.Request] = [] + cut = "x" * client_module._METADATA_CUT + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + full = "expandMetadata" in request.url.params + rows = [{"id": "o1", "traceId": "t1", "metadata": {"tool": cut + ("canary" if full else ""), "short": "s"}}] + return httpx.Response(200, json={"data": rows, "meta": {}}) + + trace = make_client(handler).get_trace("t1") + assert seen[-1].url.params["expandMetadata"] == "tool" + assert trace["observations"][0]["metadata"]["tool"].endswith("canary") + + +def test_get_trace_retries_a_throttled_read(make_client, monkeypatch): + monkeypatch.setattr(client_module.time, "sleep", lambda _s: None) + answers = iter( + [httpx.Response(429), httpx.Response(200, json={"data": [{"id": "o1", "traceId": "t1"}], "meta": {}})] + ) + + trace = make_client(lambda request: next(answers)).get_trace("t1") + assert [o["id"] for o in trace["observations"]] == ["o1"] + + +def test_get_trace_raises_on_a_client_error(make_client): + with pytest.raises(httpx.HTTPStatusError): + make_client(lambda request: httpx.Response(404)).get_trace("missing") + + +def test_trace_reads_wait_longer_than_the_client_default(make_client): + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(200, json={"data": [], "meta": {}}) + + client = make_client(handler) + client.get_trace("t1") + client.session_trace_ids("conv-1") + assert [r.extensions["timeout"]["read"] for r in seen] == [client_module._TRACE_READ_TIMEOUT] * 2 diff --git a/packages/gooddata-eval/tests/test_langfuse_source.py b/packages/gooddata-eval/tests/test_langfuse_source.py index a9db1997e..f18094c69 100644 --- a/packages/gooddata-eval/tests/test_langfuse_source.py +++ b/packages/gooddata-eval/tests/test_langfuse_source.py @@ -214,3 +214,36 @@ def test_a_blank_declaration_does_not_hide_a_real_one_behind_it(): assert _infer_test_kind({"test_kind": "agentic_search"}, "visualization", {"test_kind": "agentic_guardrail"}) == ( "agentic_search" ) + + +def test_item_from_raw_list_input_is_a_scripted_conversation(): + raw = { + "id": "lf-3", + "datasetName": "agent_obfuscation", + "input": ["My email is a@b.invalid, show revenue", "Now by quarter"], + "expectedOutput": {"canaries": [{"nonce": "n", "value": "a@b.invalid"}]}, + } + item = _item_from_raw(raw, dataset_name="agent_obfuscation", test_kind="visualization") + assert item.question == "My email is a@b.invalid, show revenue" + assert item.turns == ["My email is a@b.invalid, show revenue", "Now by quarter"] + assert item.test_kind == "agentic_obfuscation" + + +def test_item_from_raw_query_input_keeps_a_json_question_verbatim(): + # Langfuse parses a string input that is valid JSON into an object, so such a question + # travels wrapped in {"query": ...} and must come back as the exact string. + question = '{"api_key": "ak_x", "password": 90210731}' + raw = {"id": "lf-4", "input": {"query": question}, "expectedOutput": {"canaries": []}} + item = _item_from_raw(raw, dataset_name="ds", test_kind="visualization") + assert item.question == question + assert item.turns is None + + +def test_item_from_raw_rejects_an_input_that_lost_its_question(): + raw = {"id": "lf-5", "input": {"api_key": "ak_x"}, "expectedOutput": {"canaries": []}} + with pytest.raises(ValueError, match="Unsupported Langfuse item input shape"): + _item_from_raw(raw, dataset_name="ds", test_kind="visualization") + + +def test_infer_test_kind_recognises_obfuscation_canaries(): + assert _infer_test_kind({"canaries": []}, "visualization") == "agentic_obfuscation" diff --git a/packages/gooddata-eval/tests/test_models.py b/packages/gooddata-eval/tests/test_models.py index 0c09a23ef..f20d08347 100644 --- a/packages/gooddata-eval/tests/test_models.py +++ b/packages/gooddata-eval/tests/test_models.py @@ -171,3 +171,24 @@ def test_a_visualization_without_an_id_still_parses(): """The agent sometimes omits `id`; nothing scores on it, so it must not error the item.""" viz = CreatedVisualization.model_validate({"type": "bar_chart", "query": {"fields": {"m": "metric/x"}}}) assert viz.id is None + + +def test_dataset_item_accepts_a_list_of_turns_as_question(): + item = DatasetItem.model_validate( + { + "id": "1", + "dataset_name": "d", + "test_kind": "agentic_obfuscation", + "question": ["a", "b"], + "expected_output": {}, + } + ) + assert item.question == "a" + assert item.turns == ["a", "b"] + + +def test_dataset_item_single_question_has_no_turns(): + item = DatasetItem.model_validate( + {"id": "1", "dataset_name": "d", "test_kind": "agentic_guardrail", "question": "q", "expected_output": "x"} + ) + assert item.turns is None diff --git a/packages/gooddata-eval/tests/test_obfuscation_check.py b/packages/gooddata-eval/tests/test_obfuscation_check.py new file mode 100644 index 000000000..6ea488bfd --- /dev/null +++ b/packages/gooddata-eval/tests/test_obfuscation_check.py @@ -0,0 +1,197 @@ +# (C) 2026 GoodData Corporation. All rights reserved. +# SPDX-License-Identifier: LicenseRef-GoodData-Enterprise +import json + +import pytest +from gooddata_eval.core.agentic._obfuscation_check import ( + AI_ON_DATA_REDACTION, + SINK_CONVERSATION_DB, + SINK_LANGFUSE_TRACE, + SINKS, + derive_anchor, + evaluate_case, + find_canary, +) +from gooddata_eval.core.agentic.obfuscation import without_canaries + +_QUESTION = "Can you check the recent orders for the account registered to marie.canary7f3a@example.invalid?" +_EMAIL = { + "nonce": "canary7f3a", + "value": "marie.canary7f3a@example.invalid", + "class": "EMAIL", + "absent_from": [SINK_LANGFUSE_TRACE, SINK_CONVERSATION_DB], + "mask_marker_present": "[EMAIL]", +} +_ENFORCED = {"canaries": [_EMAIL], "status": "enforced"} +_ANCHOR = "Can you check the recent orders for the" +_MASKED = "Can you check the recent orders for the account registered to [EMAIL]?" + + +def _db(text: str) -> dict: + return {"conversation": {"title": None}, "items": [{"role": "user", "content": {"type": "text", "text": text}}]} + + +def _trace(text: str) -> list: + # Langfuse keeps span input as a JSON document serialised into a string. + return [{"id": "t1", "input": json.dumps({"parts": [{"text": text}]}), "observations": []}] + + +def _sinks(db_text: str, trace_text: str) -> dict: + return {SINK_CONVERSATION_DB: _db(db_text), SINK_LANGFUSE_TRACE: _trace(trace_text)} + + +def test_anchor_is_the_longest_plain_word_run_outside_the_canary(): + assert derive_anchor(_QUESTION, [_EMAIL]) == _ANCHOR + + +def test_anchor_does_not_keep_the_part_of_a_canary_left_around_its_nonce(): + escaped = { + "nonce": "canary2d8e", + "value": "john\\.canary2d8e@northwind\\-demo\\.invalid", + "unescaped_value": "john.canary2d8e@northwind-demo.invalid", + } + question = "Pouvez-vous verifier la commande passee par john\\.canary2d8e@northwind\\-demo\\.invalid la semaine" + assert derive_anchor(question, [escaped]) == "vous verifier la commande passee par" + + +def test_anchor_is_none_when_nothing_but_the_canary_is_left(): + assert derive_anchor("marie.canary7f3a@example.invalid", [_EMAIL]) is None + + +def test_masked_in_both_sinks_passes(): + verdict = evaluate_case(_ENFORCED, _sinks(_MASKED, _MASKED), [_ANCHOR]) + assert verdict.passed, verdict.failures + + +def test_canary_inside_json_serialised_trace_input_is_a_leak(): + verdict = evaluate_case(_ENFORCED, _sinks(_MASKED, _QUESTION), [_ANCHOR]) + assert [f for f in verdict.failures if f.startswith("LEAK")] == [ + f"LEAK EMAIL 'canary7f3a' in {SINK_LANGFUSE_TRACE}: $[0].input.parts[0].text ('canary7f3a')" + ] + + +def test_a_sink_without_the_anchor_is_blind_not_clean(): + verdict = evaluate_case(_ENFORCED, _sinks(_MASKED, "unrelated title-generation trace"), [_ANCHOR]) + assert not verdict.passed + assert any("blind" in f and SINK_LANGFUSE_TRACE in f for f in verdict.failures) + + +def test_a_sink_that_was_not_collected_fails(): + verdict = evaluate_case(_ENFORCED, {SINK_CONVERSATION_DB: _db(_MASKED), SINK_LANGFUSE_TRACE: None}, [_ANCHOR]) + assert any("not collected" in f for f in verdict.failures) + + +def test_ai_on_data_redaction_is_reported_instead_of_passing(): + verdict = evaluate_case(_ENFORCED, _sinks(_MASKED, AI_ON_DATA_REDACTION), [_ANCHOR]) + assert any("enableGenAiTraceRedaction" in f for f in verdict.failures) + + +def test_missing_mask_marker_fails_even_without_a_leak(): + verdict = evaluate_case(_ENFORCED, _sinks(_MASKED, "Can you check the recent orders for the account?"), [_ANCHOR]) + assert any("mask marker '[EMAIL]' missing from langfuse_trace" in f for f in verdict.failures) + + +def test_card_written_with_spaces_matches_the_compact_nonce(): + card = {"nonce": "4916338506082832", "value": "4916338506082832", "class": "CREDIT_CARD"} + assert find_canary(card, SINK_CONVERSATION_DB, {"text": "card 4916 3385-0608 2832"}) + + +def test_known_limitation_that_stops_reproducing_asks_for_a_deliberate_flip(): + person = {"nonce": "Canary8d47", "value": "Marina Canary8d47", "class": "PERSON", "present_in": list(SINKS)} + expected = { + "canaries": [person], + "status": "known_limitation", + "known_limitation": {"flip_when": "the PERSON classifier ships"}, + } + anchor = "was diagnosed in march" + masked = f"Patient [PERSON] {anchor}" + verdict = evaluate_case(expected, _sinks(masked, masked), [anchor]) + assert len(verdict.failures) == 2 + assert all("the PERSON classifier ships" in f for f in verdict.failures) + + +_DECIMAL = { + "nonce": "5621697898080164", + "value": "0.5621697898080164", + "class": "NOT_SENSITIVE", + "present_in": list(SINKS), +} +_DECIMAL_ANCHOR = "Is that higher than the average monthly" + + +def test_an_enforced_value_that_must_stay_is_reported_as_over_masked_when_masked(): + masked = f"Our margin came out as 0.[CREDIT_CARD]. {_DECIMAL_ANCHOR} margin?" + verdict = evaluate_case({"canaries": [_DECIMAL], "status": "enforced"}, _sinks(masked, masked), [_DECIMAL_ANCHOR]) + assert sorted(verdict.failures) == [ + f"OVER-MASKED NOT_SENSITIVE '5621697898080164' in {sink}: expected unchanged but it was masked (false positive)" + for sink in sorted(SINKS) + ] + + +def test_an_enforced_value_that_must_stay_passes_when_left_alone(): + kept = f"Our margin came out as 0.5621697898080164. {_DECIMAL_ANCHOR} margin?" + verdict = evaluate_case({"canaries": [_DECIMAL], "status": "enforced"}, _sinks(kept, kept), [_DECIMAL_ANCHOR]) + assert verdict.passed + + +def test_a_canary_used_as_a_json_key_is_a_leak(): + pasted = '{"marie.canary7f3a@example.invalid": true, "note": "Can you check the recent orders for the account"}' + verdict = evaluate_case(_ENFORCED, _sinks(pasted, _MASKED), [_ANCHOR]) + assert any( + f.startswith(f"LEAK EMAIL 'canary7f3a' in {SINK_CONVERSATION_DB}") and "" in f for f in verdict.failures + ) + + +def test_rejected_turn_may_export_no_trace(): + expected = {"canaries": [{**_EMAIL, "mask_marker_present": None}], "status": "enforced"} + sinks = {SINK_CONVERSATION_DB: {"conversation": {}, "items": []}, SINK_LANGFUSE_TRACE: None} + verdict = evaluate_case(expected, sinks, [], turn_rejected=True) + assert verdict.passed, verdict.failures + assert verdict.notes + + +def test_rejected_turn_still_convicts_a_trace_that_shows_the_canary(): + expected = {"canaries": [{**_EMAIL, "mask_marker_present": None}], "status": "enforced"} + sinks = {SINK_CONVERSATION_DB: {"conversation": {}, "items": []}, SINK_LANGFUSE_TRACE: _trace(_QUESTION)} + verdict = evaluate_case(expected, sinks, [], turn_rejected=True) + assert any(f.startswith("LEAK") for f in verdict.failures) + + +def test_record_only_path_is_reported_but_not_gated(): + canary = {**_EMAIL, "record_only_paths": [r"\.alertProposal\."]} + db = _db(_MASKED) + db["items"].append({"role": "assistant", "content": {"parts": [{"alertProposal": {"to": _EMAIL["value"]}}]}}) + verdict = evaluate_case({"canaries": [canary]}, {**_sinks(_MASKED, _MASKED), SINK_CONVERSATION_DB: db}, [_ANCHOR]) + assert verdict.passed, verdict.failures + assert any(n.startswith("RECORDED EMAIL") and "alertProposal" in n for n in verdict.notes) + + +def test_a_canary_outside_the_record_only_path_still_fails(): + canary = {**_EMAIL, "record_only_paths": [r"\.alertProposal\."]} + verdict = evaluate_case({"canaries": [canary]}, _sinks(_QUESTION, _MASKED), [_ANCHOR]) + assert any(f.startswith("LEAK") for f in verdict.failures) + + +def test_unknown_sink_name_is_a_fixture_error(): + expected = {"canaries": [{**_EMAIL, "absent_from": ["langfuse"]}]} + with pytest.raises(ValueError, match="unknown sink"): + evaluate_case(expected, _sinks(_MASKED, _MASKED), [_ANCHOR]) + + +def test_langfuse_comment_carries_no_canary_spelling(): + card = {"nonce": "4916338506082832", "value": "4916 3385 0608 2832", "class": "CREDIT_CARD"} + text = "LEAK CREDIT_CARD '4916338506082832' in conversation_db: $.items[0] ('4916 3385 0608 2832')" + cleaned = without_canaries(text, [card]) + assert "4916" not in cleaned + assert "" in cleaned + + +def test_a_numeric_secret_inside_decoded_json_is_found(): + password = {"nonce": "90210731", "value": "90210731", "class": "PASSWORD"} + stored = {"items": [{"content": {"text": '{"api_key": "[API_KEY]", "password": 90210731}'}}]} + assert find_canary(password, SINK_CONVERSATION_DB, stored) + + +def test_booleans_are_not_leaves(): + flag = {"nonce": "True", "value": "True", "class": "TOKEN"} + assert not find_canary(flag, SINK_CONVERSATION_DB, {"passed": True}) diff --git a/packages/gooddata-eval/tests/test_sse_client.py b/packages/gooddata-eval/tests/test_sse_client.py index ae44fd9d1..890877027 100644 --- a/packages/gooddata-eval/tests/test_sse_client.py +++ b/packages/gooddata-eval/tests/test_sse_client.py @@ -1000,3 +1000,12 @@ def handler(request): ) client.ask(item) assert captured["body"]["userContext"] == _ATTACHMENT + + +def test_parse_sse_lines_error_keeps_the_reason_gen_ai_sends_without_a_detail(): + lines = ['data: {"statusCode": 422, "reason": "DATA_OBFUSCATION_CONTENT_REJECTED"}'] + with pytest.raises(ChatError) as raised: + parse_sse_lines(lines) + assert raised.value.status_code == 422 + assert raised.value.reason == "DATA_OBFUSCATION_CONTENT_REJECTED" + assert "DATA_OBFUSCATION_CONTENT_REJECTED" in str(raised.value) diff --git a/packages/gooddata-eval/tests/test_trace_linker.py b/packages/gooddata-eval/tests/test_trace_linker.py index 2801e349f..9e30b38f6 100644 --- a/packages/gooddata-eval/tests/test_trace_linker.py +++ b/packages/gooddata-eval/tests/test_trace_linker.py @@ -120,6 +120,7 @@ def test_run_trace_link_inline_runs_the_task_on_the_calling_thread(): ("kda_skill", "evaluate_agentic_kda_skill"), ("conversation", "evaluate_agentic_conversation"), ("what_if", "evaluate_agentic_what_if"), + ("obfuscation", "evaluate_agentic_obfuscation"), ]