Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
333 changes: 333 additions & 0 deletions doc/code/analytics/0_attack_results.md

Large diffs are not rendered by default.

7 changes: 6 additions & 1 deletion doc/code/framework.md
Original file line number Diff line number Diff line change
Expand Up @@ -334,9 +334,13 @@ The below talks about responsibilities of most modules in the PyRIT library
- This is where cross-run analysis belongs: e.g. "which attack performed best for this objective?", "how often did a technique succeed?", or "which responses match known content?".
- **Does not own**: live, in-attack decisions — any decision made *during* an attack is a scorer's job. Analytics only operates on stored results, after the fact.
- Today it includes `ConversationAnalytics` (inspecting conversation history), `analyze_results` / `AttackStats` (aggregating outcomes across techniques), and text-matching strategies (`ExactTextMatching`, `ApproximateTextMatching`).
- `compute_scenario_statistics` calculates scenario success statistics. It owns execution-unit identity (atomic attack, technique configuration, and seed group), latest-attempt selection, counts, denominators, and rounding. SDK callers, the GUI backend's run detail and progress views, and the console, JSON, and HTML reports all present its results (`ScenarioExecutionStatistics`, `ScenarioExecutionUnit`, and `ScenarioProgressCounts` in `pyrit.models`) instead of calculating their own. The one exception is the GUI run-history list, which aggregates the same statistics in SQL (`MemoryInterface._build_scenario_history_aggregate_statement`) so it can page over many runs; `tests/unit/analytics/test_scenario_statistics_parity.py` keeps the two implementations in agreement.
- `compute_scenario_statistics` owns execution-unit identity (atomic attack, technique configuration, and seed group), latest-attempt selection, and historical retry/error counts. It passes selected outcomes to the shared outcome calculator. SDK callers, the GUI backend's run detail and progress views, and the console, JSON, and HTML reports present its results (`ScenarioExecutionStatistics`, `ScenarioExecutionUnit`, and `ScenarioProgressCounts` in `pyrit.models`). The GUI run-history list aggregates counts in SQL (`MemoryInterface._build_scenario_history_aggregate_statement`) so it can page over many runs, then uses the shared percentage helper; `tests/unit/analytics/test_scenario_statistics_parity.py` keeps the paths in agreement.
- Scenario attempts are ordered by timestamp, then their canonical lowercase UUID string. SQL Server uses this string order rather than its native UUID order for history ranking and progress pagination. Explicit seed attribution wins; an objective alone matches a planned seed group only when that match is unique. Legacy runs with identifier-only seed identities are recounted with the shared analytics.
- Shared analytics contracts (filters, dimensions, typed values, reports, facets, result pages, and `AttackStats`) live in `pyrit.models.analytics`. They validate data without querying memory or calculating statistics. `AttackResultSelection` defines selection modes without changing existing callers.
- `compute_outcome_statistics` owns both success-rate denominators, totals, and shares. `success_rate_decided` (the explicit alias of `success_rate`) divides by successes plus failures; `success_rate_all` divides by all outcomes. `combine_outcome_statistics` revalidates and combines disjoint counts, never averages rates. Attack and scenario analytics share the same `OutcomeStatistics` model and calculations after independently selecting their populations. Models reject contradictory supplied totals/rates; they do not select populations or replace the supplied numbers.
- `analyze_results` and `compute_technique_stats_async` remain maintained APIs. Their `include_outcome_statistics=True` option exposes the shared statistics directly; default `AttackStats` results retain their six-field shape. Existing scenario percentage fields also retain their defaults. Only the already-deprecated sync memory wrappers and `ScenarioResult.objective_achieved_rate` are scheduled for removal in 1.4.0.
- [`AttackResultAnalytics`](./analytics/0_attack_results.md) provides async saved-result reports, lightweight result pages, and facet lookups. It supplies both shared success rates, display labels and exact additional drill-down predicates, and counts every saved result ID rather than latest scenario execution units. Memory owns cohort SQL and consistent projections; the SDK owns their interpretation.
- Analytics facades sharing a memory instance share one loop-bound, bounded report/quick-query controller. Close that shared SDK lifetime before replacing or disposing memory; cancellation does not free a running query's slot before its actual session cleanup. SDK lifecycle is independent of backend runtime integration.
- Filter-bound cursor and label-normalization helpers live in `pyrit.common.pagination`. The backend pagination module retains compatibility exports, including History's invalid-cursor first-page fallback.

## Auth
Expand All @@ -362,6 +366,7 @@ The below talks about responsibilities of most modules in the PyRIT library
- Components should access memory through `CentralMemory` rather than passing state directly between each other.
- Memory backends are swappable too (e.g. SQLite or Azure SQL) without changing the components that use them.
- Memory loads and locks observation evidence for model-owned validation, and owns atomic writes and reference cleanup.
- Analytics readers return saved-outcome counts and lightweight metadata, or bounded typed profiles with a complete SQL fallback. They own native async sessions, consistent read views, and `QueryControl` database deadlines/interruption, not rate calculations or display labels.
- **Does not own**: business logic or decisions. Memory stores and retrieves state; it doesn't decide what to send, how to score, or when to branch — components do that and persist results here.

## [Models](../contributing/11_memory_models)
Expand Down
1 change: 1 addition & 0 deletions doc/myst.yml
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,7 @@ project:
- file: code/registry/1_class_registry.ipynb
- file: code/registry/2_instance_registry.ipynb
- file: code/output/0_output.ipynb
- file: code/analytics/0_attack_results.md
- file: api/index.md
children:
- file: api/pyrit_analytics.md
Expand Down
7 changes: 7 additions & 0 deletions pyrit/analytics/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,21 +9,28 @@
from pyrit.common.lazy_imports import get_lazy_dir, resolve_lazy_export

if TYPE_CHECKING:
from pyrit.analytics.attack_result_analytics import AttackResultAnalytics
from pyrit.analytics.conversation_analytics import ConversationAnalytics
from pyrit.analytics.outcome_statistics import combine_outcome_statistics, compute_outcome_statistics
from pyrit.analytics.result_analysis import (
AttackStats,
analyze_results,
get_cached_results_for_technique,
get_cached_results_for_technique_async,
)
from pyrit.analytics.scenario_statistics import compute_scenario_statistics
from pyrit.analytics.technique_analysis import compute_technique_stats_async
from pyrit.analytics.text_matching import ApproximateTextMatching, ExactTextMatching, TextMatching

_LAZY_EXPORTS: dict[str, str | tuple[str, str | None]] = {
"analyze_results": "pyrit.analytics.result_analysis",
"ApproximateTextMatching": "pyrit.analytics.text_matching",
"AttackResultAnalytics": "pyrit.analytics.attack_result_analytics",
"AttackStats": "pyrit.analytics.result_analysis",
"combine_outcome_statistics": "pyrit.analytics.outcome_statistics",
"compute_outcome_statistics": "pyrit.analytics.outcome_statistics",
"compute_scenario_statistics": "pyrit.analytics.scenario_statistics",
"compute_technique_stats_async": "pyrit.analytics.technique_analysis",
"ConversationAnalytics": "pyrit.analytics.conversation_analytics",
"ExactTextMatching": "pyrit.analytics.text_matching",
"get_cached_results_for_technique": "pyrit.analytics.result_analysis",
Expand Down
278 changes: 278 additions & 0 deletions pyrit/analytics/_execution.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,278 @@
# Copyright (c) Microsoft Corporation.
# Licensed under the MIT license.

"""Bound native async analytics work through its actual resource cleanup."""

from __future__ import annotations

import asyncio
import logging
from collections import deque
from dataclasses import dataclass, field
from sys import float_info
from time import monotonic
from typing import TYPE_CHECKING, TypeVar

from pyrit.common.task_utils import gather_with_cleanup_async
from pyrit.exceptions.analytics_exception import AnalyticsBusyException, AnalyticsTimeoutException
from pyrit.memory.query_control import QueryControl

if TYPE_CHECKING:
from collections.abc import Awaitable, Callable, Coroutine

logger = logging.getLogger(__name__)
T = TypeVar("T")


@dataclass(eq=False)
class _Waiter:
ready: asyncio.Future[None]
deadline: float


@dataclass(eq=False)
class _Operation:
"""Keep the task strongly owned even after its requesting coroutine has left."""

control: QueryControl
finished: asyncio.Future[None]
abandoned: bool = False
task: Awaitable[object] | None = None


@dataclass
class _Lane:
limit: int
timeout: float
active: int = 0
queued: deque[_Waiter] = field(default_factory=deque)
running: set[_Operation] = field(default_factory=set)


class AnalyticsExecution:
"""
Reserve independent report and quick-query capacity on one event loop.

Admission is bounded before creating an operation task. A response deadline or
caller cancellation signals ``QueryControl`` but does not cancel that task:
native drivers and session cleanup must finish before its slot is released.
There is no coalescing, result cache, thread pool, or SQL timeout policy here.
The reader owns database interruption and session cleanup.
"""

def __init__(
self,
*,
report_workers: int = 5,
quick_workers: int = 2,
max_queue: int = 10,
queue_timeout: float = 1.0,
report_timeout: float = 5.0,
quick_timeout: float = 1.0,
) -> None:
"""
Bind a controller to the running loop without starting database work.

Args:
report_workers (int): Concurrent report slots, including cleanup.
quick_workers (int): Separate slots shared by result pages and facets.
max_queue (int): Maximum waiting requests per lane, excluding active slots.
queue_timeout (float): Maximum admission wait in seconds.
report_timeout (float): Report execution budget after admission, in seconds.
quick_timeout (float): Result-page/facet execution budget, in seconds.

Raises:
ValueError: If a count is not a positive integer or a timeout is not positive and finite.
RuntimeError: If construction occurs outside a running event loop.
"""
for name, value in (
("report_workers", report_workers),
("quick_workers", quick_workers),
("max_queue", max_queue),
):
if type(value) is not int or value <= 0:
raise ValueError(f"{name} must be a positive integer.")
for name, timeout in (
("queue_timeout", queue_timeout),
("report_timeout", report_timeout),
("quick_timeout", quick_timeout),
):
if isinstance(timeout, bool) or not isinstance(timeout, (int, float)) or not 0 < timeout <= float_info.max:
raise ValueError(f"{name} must be positive and finite.")
self._loop = asyncio.get_running_loop()
self._lanes = {
True: _Lane(limit=report_workers, timeout=report_timeout),
False: _Lane(limit=quick_workers, timeout=quick_timeout),
}
self._max_queue = max_queue
self._queue_timeout = queue_timeout
self._closing = False
self._closed = False
self._close_task: asyncio.Task[None] | None = None

@property
def is_closed(self) -> bool:
"""Whether shutdown has drained all started operations, not just their callers."""
return self._closed

async def run_async(self, *, report: bool, task: Callable[[QueryControl], Awaitable[T]]) -> T:
"""
Admit one operation and await its response without blocking the owning loop.

Args:
report (bool): Use the report lane rather than the reserved quick lane.
task (Callable[[QueryControl], Awaitable[T]]): Native async work that honors
the shared control and does not return before releasing its resources.

Returns:
T: This operation's result. Separate calls never share response objects.

Raises:
AnalyticsBusyException: If the lane is full, admission expires, or shutdown has begun.
AnalyticsTimeoutException: If execution exceeds its budget.
CancelledError: If the caller leaves; running work retains its slot until cleanup finishes.
RuntimeError: If used from a different event loop.
"""
self._check_loop()
lane = self._lanes[report]
await self._admit_async(lane)
work = _Operation(
control=QueryControl(deadline=monotonic() + lane.timeout),
finished=self._loop.create_future(),
)
lane.running.add(work)
try:
operation = self._create_task(self._execute_async(task=task, control=work.control))
except BaseException:
self._finish(lane=lane, work=work)
raise
work.task = operation
operation.add_done_callback(lambda completed: self._complete(lane=lane, work=work, task=completed))
try:
done, _ = await asyncio.wait({operation}, timeout=work.control.remaining)
if not done:
self._abandon(work=work, task=operation)
raise AnalyticsTimeoutException
return operation.result()
except asyncio.CancelledError:
self._abandon(work=work, task=operation)
raise

async def close_async(self) -> None:
"""
Reject admission, signal cancellation, and drain actual work on the owning loop.

This is idempotent and terminal. It can outlast response deadlines because
returning early would allow a replacement controller to overlap database
operations still cleaning up. Cancelling this await, even repeatedly, is
propagated only after draining. The memory backend itself is not disposed.
If scheduling the drain fails, admission remains closed and a later
``close_async`` call can retry the drain.
"""
self._check_loop()
if self._close_task is None:
self._closing = True
for lane in self._lanes.values():
while lane.queued:
waiter = lane.queued.popleft()
if not waiter.ready.done():
waiter.ready.set_exception(AnalyticsBusyException())
for work in lane.running:
work.control.cancel()
self._close_task = self._create_task(self._drain_async())
cancellation: asyncio.CancelledError | None = None
while not self._close_task.done():
try:
await asyncio.shield(self._close_task)
except asyncio.CancelledError as error:
cancellation = error
self._close_task.result()
if cancellation is not None:
raise cancellation

def _check_loop(self) -> None:
if asyncio.get_running_loop() is not self._loop:
raise RuntimeError("Analytics must be used and closed on its owning event loop.")

def _create_task(self, coroutine: Coroutine[object, object, T]) -> asyncio.Task[T]:
try:
return self._loop.create_task(coroutine)
except BaseException:
coroutine.close()
raise

async def _admit_async(self, lane: _Lane) -> None:
if self._closing:
raise AnalyticsBusyException
if lane.active < lane.limit:
lane.active += 1
return
if len(lane.queued) >= self._max_queue:
raise AnalyticsBusyException
waiter = _Waiter(ready=self._loop.create_future(), deadline=monotonic() + self._queue_timeout)
lane.queued.append(waiter)
try:
try:
async with asyncio.timeout(self._queue_timeout):
await waiter.ready
except TimeoutError as error:
raise AnalyticsBusyException from error
if self._closing or monotonic() >= waiter.deadline:
raise AnalyticsBusyException
except BaseException:
if waiter in lane.queued:
lane.queued.remove(waiter)
elif not waiter.ready.cancelled() and waiter.ready.exception() is None:
self._release(lane)
raise

def _release(self, lane: _Lane) -> None:
lane.active -= 1
while lane.queued and not self._closing:
waiter = lane.queued.popleft()
if waiter.ready.done():
continue
if monotonic() >= waiter.deadline:
waiter.ready.set_exception(AnalyticsBusyException())
continue
lane.active += 1
waiter.ready.set_result(None)
break

def _complete(self, *, lane: _Lane, work: _Operation, task: asyncio.Task[T]) -> None:
if not task.cancelled():
task.exception()
self._finish(lane=lane, work=work)
if work.abandoned:
self._log_abandoned_error(task)

def _finish(self, *, lane: _Lane, work: _Operation) -> None:
lane.running.remove(work)
work.task = None
self._release(lane)
work.finished.set_result(None)

def _abandon(self, *, work: _Operation, task: asyncio.Task[T]) -> None:
work.abandoned = True
work.control.cancel()
# Completion can precede cancellation delivery, after its callback consumed the exception.
if work.finished.done():
self._log_abandoned_error(task)

async def _drain_async(self) -> None:
await gather_with_cleanup_async(work.finished for lane in self._lanes.values() for work in lane.running)
self._closed = True

@staticmethod
async def _execute_async(*, task: Callable[[QueryControl], Awaitable[T]], control: QueryControl) -> T:
control.check()
result = await task(control)
control.check()
return result

@staticmethod
def _log_abandoned_error(task: asyncio.Task[T]) -> None:
if not task.cancelled():
error = task.exception()
if error is not None and not isinstance(error, AnalyticsTimeoutException):
logger.error("Analytics operation failed after its caller left.", exc_info=error)
Loading