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
55 changes: 55 additions & 0 deletions docs/design/conversation-compaction/conversation-compaction.md
Original file line number Diff line number Diff line change
Expand Up @@ -315,6 +315,7 @@ Add `compaction` field to the root `Configuration` class.
| `src/models/config.py` | Add `CompactionConfiguration` (near `ConversationHistoryConfiguration`) |
| `src/configuration.py` | Add `compaction_configuration` property to `AppConfig` singleton |
| `src/utils/conversation_compaction.py` | New module: `apply_compaction()` / `apply_compaction_blocking()`, `needs_compaction_path()`, marker helpers, per-conversation lock |
| `src/utils/pending_turn.py` | `PendingTurn`, the owner of every turn the endpoints store themselves: decides whether the turn is ours, stores it against the input as it arrived, and stores it once — LCORE-3908 |
| `src/models/common/responses/responses_api_params.py` | `omit_conversation` flag — drops the `conversation` parameter from the request body in compacted mode |
| `src/app/endpoints/query.py` | Call `apply_compaction_blocking()` after preparing params; store the turn in compacted mode |
| `src/app/endpoints/streaming_query.py` | Compaction-aware SSE path that emits the `compaction` event before summarizing (R12) |
Expand Down Expand Up @@ -347,6 +348,60 @@ When compaction is active, the endpoint builds explicit input, the
`conversation` parameter is omitted (via `ResponsesApiParams.omit_conversation`),
and the completed turn is appended to the conversation items afterward.

That append has one owner, `PendingTurn` in `src/utils/pending_turn.py`
(LCORE-3908). An endpoint creates it from the parameters the request is sent
with and the `original_input` of the `CompactionResult`, and reports how the
turn ended: `store_completed`, `store_blocked` or `store_interrupted`.
The first report settles the turn and every later one does nothing, so a turn is
stored once even when several paths of a request want to store it (the end of a
stream, the cancellation handler, the interrupt callback). `PendingTurn` also
covers the other turns OGX does not store: a request a shield blocked on
`/v1/responses`, an interrupted stream, a continuation from
`previous_response_id`.

A compacted request that loses its turn through a change in the code fails.
Building a `PendingTurn` for compacted parameters without the original input
raises `ValueError`, and `ensure_settled()` raises `TurnNotStoredError` when a
compacted request reaches the end of its handler and nobody tried to store its
turn. `/v1/query` runs inside the `pending_turn()` scope, which makes that check
when it is left; the streaming paths, `/v1/responses` and A2A make it after
their write. Before this, a lost write was silent: the conversation stopped
growing and nothing failed (LCORE-3883).

The check does not cover three endings, which store nothing and are rows of the
table below: a stream the client stops reading, a `/v1/responses` stream without
a final response, and a failed write where the failure is logged.

The shield capabilities (question validity, Granite Guardian) are the one writer
outside the owner. They store the turn they rejected from inside the agent run,
and only when the model was handed the conversation, so never in compacted mode.

What is stored, per endpoint and per way a turn can end. "OGX" means the
`conversation` parameter was sent and OGX stores the turn itself. The last
column says what a failed write does to the request.

| Endpoint | Turn ended | Not compacted | Compacted | Failed write |
|---|---|---|---|---|
| `/v1/query` | completed | OGX | stored | request fails |
| | model call failed | not stored | not stored | |
| | run did not finish with success | OGX | not stored | |
| `/v1/streaming_query` | completed, also when the run did not finish with success | OGX | stored when the stream ends | logged |
| | interrupted by the client | stored, with the answer so far | stored, with the answer so far | logged |
| | client stopped reading | OGX | not stored | |
| `/v1/responses` | completed, incomplete or failed | OGX; stored when continuing from `previous_response_id` | stored | request fails; a stream ends before `[DONE]` |
| | blocked by a shield | stored | stored | request fails |
| | stream without a final response | not stored | not stored | |
| | `store: false` | not stored | never compacted | |
| A2A | completed | OGX | stored | logged |
| | agent run failed | not stored | not stored | |

`tests/integration/endpoints/test_turn_persistence.py` pins the table. It runs
the real handlers and compares the conversation item by item after the request.
The compacted column is covered row by row; the other column for the rows where
lightspeed-stack stores the turn, and for a completed turn on each endpoint. The
last column is covered for a completed turn on each endpoint and for an
interrupted stream.

## Fetching conversation history

Use the same pattern as `conversations_v1.py:240-246`:
Expand Down
2 changes: 1 addition & 1 deletion docs/devel_doc/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -391,7 +391,7 @@ The compaction system is split into two layers:
- Marker persistence (`[lightspeed:compaction-summary]` sentinel in conversation items)
- `CompactionStartedEvent` emission for streaming progress indicators
- `apply_compaction()` (async generator) — Main entry point used by all endpoints
- `store_compacted_turn()` — Appends user query + LLM output when in compacted mode
- `PendingTurn` (`utils/pending_turn.py`) — Owns the turns the endpoints store themselves: appends user query + LLM output when in compacted mode, once per request

**Data Flow:**

Expand Down
81 changes: 44 additions & 37 deletions src/app/endpoints/a2a.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@
)
from a2a.utils import new_agent_text_message, new_task
from fastapi import APIRouter, Depends, HTTPException, Request, status
from ogx_client import ApiException
from ogx_client import ApiException, AsyncOgxClient
from opentelemetry import trace
from pydantic_ai import AgentRunResultEvent
from pydantic_ai.exceptions import AgentRunError
Expand Down Expand Up @@ -63,11 +63,7 @@
from models.common.responses.responses_api_params import ResponsesApiParams
from models.config import Action
from utils.agents.error_handler import map_agent_inference_error
from utils.conversation_compaction import (
CompactionResult,
apply_compaction_blocking,
store_compacted_turn,
)
from utils.conversation_compaction import apply_compaction_blocking
from utils.mcp.mcp_headers import McpHeaders, mcp_headers_dependency
from utils.otel_tracing import (
SpanAttributes,
Expand All @@ -76,6 +72,7 @@
anonymize_value,
set_span_attributes,
)
from utils.pending_turn import PendingTurn
from utils.pydantic_ai_helpers import build_agent, captured_output_items
from utils.query import extract_provider_and_model_from_model_id
from utils.responses import prepare_responses_params
Expand Down Expand Up @@ -213,10 +210,39 @@ def _record_execution_span(
span.set_attribute(SpanAttributes.INFERENCE_TIME, inference_time)


async def _persist_compacted_a2a_turn(
client: Any,
async def _compact_a2a_request(
client: AsyncOgxClient,
responses_params: ResponsesApiParams,
compaction: CompactionResult,
) -> tuple[ResponsesApiParams, PendingTurn]:
"""Compact the conversation of an A2A request if it nears the context window.

A2A is not a browser SSE stream, so no progress event is emitted; the
blocking variant summarizes inline before the call. No conversation cache
is passed: the A2A executor has no resolved user_id for the
(user_id, conversation_id) cache key, so A2A runs in marker-only mode
(additive summaries, no persisted fold).

Parameters:
client: OGX client.
responses_params: Prepared Responses API parameters.

Returns:
The parameters to send, and the pending turn of the request, which
stores the turn when OGX does not (LCORE-3908).
"""
compaction = await apply_compaction_blocking(
client,
responses_params,
configuration.inference,
configuration.compaction,
)
return compaction.params, PendingTurn.for_request(
client, compaction.params, compaction.original_input
)


async def _persist_compacted_a2a_turn(
turn: PendingTurn,
agent: Any,
task_id: str,
) -> None:
Expand All @@ -228,22 +254,13 @@ async def _persist_compacted_a2a_turn(
context.

Parameters:
client: OGX client used to write the conversation items.
responses_params: Prepared Responses API parameters.
compaction: Outcome of applying compaction. Nothing is written unless
the request was served in compacted mode.
turn: The pending turn of the request. Nothing is written when OGX
stores the turn itself.
agent: The pydantic-ai agent whose model captured the output items.
task_id: A2A task identifier, used for error reporting.
"""
if not compaction.compacted or compaction.original_input is None:
return
try:
await store_compacted_turn(
client,
responses_params.conversation,
compaction.original_input,
captured_output_items(agent),
)
await turn.store_completed(captured_output_items(agent))
except Exception: # pylint: disable=broad-except
# The caller already has its answer; the cost of the failure is that the
# next turn in this context loses this one.
Expand Down Expand Up @@ -479,19 +496,9 @@ async def _process_task_streaming( # pylint: disable=too-many-locals,too-many-s
store=True,
request_headers=self.request_headers,
)
# Compact the conversation if it is approaching the context window
# limit. A2A is not a browser SSE stream, so no progress event is
# emitted; the blocking variant summarizes inline before the call.
# No conversation cache is passed: the A2A executor has no resolved
# user_id for the (user_id, conversation_id) cache key, so A2A runs
# in marker-only mode (additive summaries, no persisted fold).
compaction = await apply_compaction_blocking(
client,
responses_params,
configuration.inference,
configuration.compaction,
responses_params, turn = await _compact_a2a_request(
client, responses_params
)
responses_params = compaction.params

_record_model_span(span, responses_params.model)
agent = build_agent(
Expand Down Expand Up @@ -582,15 +589,15 @@ async def _process_task_streaming( # pylint: disable=too-many-locals,too-many-s
)
return

await _persist_compacted_a2a_turn(
client, responses_params, compaction, agent, task_id
)
await _persist_compacted_a2a_turn(turn, agent, task_id)
turn.ensure_settled()

_record_execution_span(
span,
self._tool_call_names,
self._run_result,
compaction.compacted,
# compacted mode: the conversation parameter was not sent
responses_params.omit_conversation,
(datetime.now(UTC) - started_at).total_seconds(),
)

Expand Down
25 changes: 15 additions & 10 deletions src/app/endpoints/query.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
root_span_turn_attributes,
set_span_attributes,
)
from utils.pending_turn import pending_turn
from utils.query import (
consume_query_tokens,
store_query_results,
Expand Down Expand Up @@ -259,16 +260,20 @@ async def _handle_query_with_tracing(
if a.content_type in IMAGE_CONTENT_TYPES
] or None

# Retrieve response using Responses API
turn_summary = await retrieve_agent_response(
client,
responses_params,
endpoint_path,
compaction.original_input if compaction.compacted else None,
shield_ids=query_request.shield_ids,
no_tools=bool(query_request.no_tools),
image_attachments=image_attachments,
)
# Retrieve response using Responses API. In compacted mode OGX does not
# store the turn; the scope fails the request if nobody tried to store it.
async with pending_turn(
client, responses_params, compaction.original_input
) as turn:
turn_summary = await retrieve_agent_response(
client,
responses_params,
endpoint_path,
turn,
shield_ids=query_request.shield_ids,
no_tools=bool(query_request.no_tools),
image_attachments=image_attachments,
)

# Combine inline RAG results (BYOK + Solr) with tool-based RAG results for the transcript
rag_chunks = inline_rag_context.rag_chunks
Expand Down
Loading
Loading