Skip to content

LCORE-3908: give the turns we store ourselves one owner - #2796

Open
max-svistunov wants to merge 1 commit into
lightspeed-core:mainfrom
max-svistunov:lcore-3908-compacted-turn-persistence-owner
Open

max-svistunov wants to merge 1 commit into
lightspeed-core:mainfrom
max-svistunov:lcore-3908-compacted-turn-persistence-owner

Conversation

@max-svistunov

@max-svistunov max-svistunov commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Description

LCORE-3908. OGX appends a turn to the conversation only when it is handed the conversation parameter and runs the inference. In compacted mode that parameter is dropped, so lightspeed-stack appends the turn itself. Every endpoint made that write on its own and nothing enforced it: two unrelated cleanups removed the calls, and for six weeks conversations stopped growing after a compaction without anything failing (LCORE-3883).

Changes since the approval (2026-09-30). The approved commit conflicted with main. It is rebased, and it differs from what was approved only where main changed underneath it. Of its 25 files, 11 have the same changed lines as approved, 2 are no longer touched, and 12 differ as listed here.

  1. main removed the shield moderation that ran before the agent on /v1/query and /v1/streaming_query (LCORE-4338: Remove obsolete pre-agent shield moderation from query and streaming. #2780). The two branches that stored a blocked turn there are gone, and with them what existed only for them:
    • src/utils/agents/query.py, src/utils/agents/streaming.py: the edits to those two branches are gone, with the removal of the import they used (main removed it). retrieve_agent_response is called without moderation_result, as on main.
    • retrieve_agent_response_generator loses the turn parameter and is as on main. src/app/endpoints/streaming_query.py no longer creates a PendingTurn for the not-compacted path: it was there to share one owner between the refusal written before the stream and the stream itself. generate_agent_response creates the owner from the parameters, as it already did when none was passed.
    • PendingTurn.drop() is removed with its only caller, the /v1/query branch, and so is the logger it used. The docstrings, a comment in src/app/endpoints/query.py and the message of TurnNotStoredError no longer mention a dropped turn.
    • tests/integration/endpoints/test_turn_persistence.py: the five tests of that path are deleted (six cases; 50 tests become 44).
    • tests/unit/utils/test_pending_turn.py: test_a_dropped_turn_is_not_stored_later and the dropped case of test_settled_turn_passes_the_check are deleted, and three docstrings are reworded. tests/unit/utils/agents/: the edits to the four tests of the blocked branches are gone, because main deleted those tests.
    • tests/unit/app/endpoints/test_streaming_query.py is no longer touched. Its only changes were mock fields for the owner the handler created.
    • Design document: the table loses the rows "blocked by a shield" of /v1/query and /v1/streaming_query and the row "blocked by a shield, then interrupted", with the paragraph about that row. The text around the table no longer mentions drop, the sentence on what the test covers no longer names a blocked stream, and "a request a shield blocked" now reads "a request a shield blocked on /v1/responses". docs/devel_doc/query_endpoint.md is no longer touched: its one sentence described the removed path.
    • The one case where behaviour changed on purpose was on that path. That paragraph is gone from this description: what is stored, and when, does not change.
  2. main records on the A2A span whether the turn was compacted, and reads it from compaction.compacted (LCORE-3755: Emit additional Otel attributes for evaluation #2818). This PR moves the compaction call into a helper, so src/app/endpoints/a2a.py passes responses_params.omit_conversation there. omit_conversation is set in one place, together with compacted=True.
  3. main added replace_last_assistant_message, with which Granite Guardian replaces the answer OGX stored (LCORE-3391: added output granite guardian #2761). It writes to a conversation, so the scan of src in tests/unit/utils/test_pending_turn.py lists it, with the three modules that mention it.
  4. tests/unit/app/endpoints/test_a2a.py: one two-line hunk is gone, because main has the same two lines now (LCORE-3755: Emit additional Otel attributes for evaluation #2818).
  5. The commit message is rewritten to match.

src/app/endpoints/responses.py is one of the 11: main changed that file too (#2780, #2781, #2818), and the same lines sit in the new code.

The owner. The write now has one owner, PendingTurn in src/utils/pending_turn.py. 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 owner decides whether the turn is ours, which input is stored, and it stores a turn once: the first report settles the turn, every later one does nothing.

It covers all seven places that stored a turn, because they are the same duty, a turn OGX does not store: a conversation in compacted mode, a request a shield blocked on /v1/responses, a stream the client interrupted, a continuation from previous_response_id. store_compacted_turn is removed. No endpoint appends a turn except through the owner.

A compacted request that loses its turn through a change in the code fails.

  • Creating the owner for compacted parameters without the original input raises ValueError. That is how a caller that lost the hand-over looks.
  • ensure_settled() raises TurnNotStoredError when a compacted request reaches the end of its handler and nobody tried to store its turn. /v1/query runs in the pending_turn() scope, which makes the check when it is left. The streaming paths, /v1/responses and A2A make it after their write, in the function that calls the write and not in the one that performs it.
  • TurnNotStoredError is not a RuntimeError, because the endpoints report that one as an inference failure.

What is kept. Every condition of the old call sites: where the write happens, what is stored for a failed or incomplete turn, what a failed write does to the request, and the order relative to quota consumption. The design document now has the table. The guard that lets only one of stream end, cancellation handler and interrupt callback finish a turn is unchanged. What is stored, and when, does not change.

What this does not do.

  • The check is a call, so a rewrite that removes the scope or the check together with the write is caught by the tests only, not by the structure.
  • The helpers take the owner as an optional parameter. Leaving it out raises ValueError in compacted mode; outside compacted mode the helper creates an owner of its own, which is what each call site did before.
  • The shield capabilities (question_validity, granite_guardian) keep storing the turn they rejected themselves. They write only when the model was handed the conversation, so never in compacted mode. /v1/query and /v1/streaming_query have no blocked turn of their own to store: the shield moderation that ran before the agent on those two endpoints is gone from main (LCORE-4338: Remove obsolete pre-agent shield moderation from query and streaming. #2780).
  • Three endings store nothing in compacted mode, as before, and the check does not cover them: a stream the client stops reading, a /v1/responses stream without a final response, a failed write that is logged. Tests record all three.
File Change
src/utils/pending_turn.py new: PendingTurn, pending_turn(), TurnNotStoredError
src/utils/agents/query.py, src/app/endpoints/query.py the owner in place of original_input; the scope
src/utils/agents/streaming.py, src/utils/stream_interrupts.py, src/app/endpoints/streaming_query.py the owner in place of original_input; the check
src/app/endpoints/responses.py both helpers store through the owner; the handlers check
src/app/endpoints/a2a.py the owner; compaction moved to a helper; the span takes the compacted flag from the parameters
src/utils/conversation_compaction.py store_compacted_turn removed
tests/integration/endpoints/test_turn_persistence.py new: the regression guard
tests/unit/utils/test_pending_turn.py new: the owner, the hand-over, the writers
tests/unit/... existing tests follow the new parameter
tests/integration/endpoints/_compaction_helpers.py, test_compaction_a2a.py the A2A helpers moved to the shared module
docs/design/conversation-compaction/conversation-compaction.md, docs/devel_doc/ARCHITECTURE.md the owner; the table

Relation to other PRs. #2791 (LCORE-4219) and #2794 (LCORE-3909) are merged, and this commit is rebased onto them. The owner lives in a module of its own because conversation_compaction.py is at 890 lines with this change, and pylint stops at 1000.

Type of change

  • Refactor
  • New feature
  • Bug fix
  • CVE fix
  • Optimization
  • Documentation Update
  • Configuration Update
  • Bump-up service version
  • Bump-up dependent library [pyproject.toml + uv.lock]
  • Bump-up dependent library [requirements.*.txt for Konflux]
  • Bump-up library or tool used for development (does not change the final image)
  • CI configuration change
  • Konflux configuration change
  • Unit tests improvement
  • Integration tests improvement
  • End to end tests improvement
  • Benchmarks improvement

Tools used to create PR

Identify any AI code assistants used in this PR (for transparency and review context)

  • Assisted-by: Claude Opus 4.8
  • Generated by: Claude Opus 4.8

Related Tickets & Documents

  • Related Issue # LCORE-3908
  • Closes # LCORE-3908

Checklist before requesting a review

  • I have performed a self-review of my code.
  • PR has passed all pre-merge test jobs.
  • If it is a core feature, I have added thorough tests.

Testing

Steps 1, 4 and 5 were run on the commit as it is now. The mutation runs of step 2 were made on the rebased commit before its last corrections, which touched no file in src/: one sentence of the design document, a test docstring, the entry for replace_last_assistant_message in the scan of src, and tests/unit/app/endpoints/test_streaming_query.py, which is back as on main. Step 3 was run on the rebased commit as it is now, on 2026-10-07, in library mode.

  1. Run the regression guard. Every test runs the real handler and compares the conversation item by item afterwards:
uv run pytest tests/integration/endpoints/test_turn_persistence.py -q   # 44 passed

The same file against the src/ of main shows that behaviour is kept (with the test helpers it imports, and src/utils/pending_turn.py added so that it can be collected). 39 pass. The five that fail take the write out of an endpoint, which main does not notice:

TestTurnNobodyStored::test_query_fails
TestTurnNobodyStored::test_streaming_query_fails
TestTurnNobodyStored::test_a2a_fails_the_task
TestTurnNobodyStored::test_responses_request_fails[blocking]
TestTurnNobodyStored::test_responses_request_fails[streaming]
  1. Check that the guard bites. Each line is one kind of change to src/, made at one site at a time (30 changes in all). Each change is followed by the tests of the paths the PR touches: tests/integration/endpoints/test_turn_persistence.py, tests/unit/utils/test_pending_turn.py, tests/unit/utils/agents, tests/unit/utils/test_stream_interrupts.py, tests/unit/app/endpoints/test_a2a.py and tests/unit/app/endpoints/test_responses.py (298 tests).
Change to src/ Sites Result
write removed: /v1/query, /v1/streaming_query, the interrupt path, A2A, /v1/responses (completed and blocked, blocking and streaming) 8 a test fails for each
check removed: the scope, /v1/streaming_query, A2A, /v1/responses (blocking, streaming), the ValueError for a missing original input 6 a test fails for each
check removed after the write of a blocked /v1/responses stream 1 no test fails
owner writes twice; a second owner writes the same turn 2 a test fails for each
owner ignores that the turn is settled 2 a test fails for each
interrupt guard ignored in the cancellation handler, in the interrupt callback 2 a test fails for each
interrupt guard ignored at the end of the stream 1 no test fails
explicit rewrite stored in place of the input 1 a test fails
a turn OGX stores is stored by us as well 2 a test fails for each
owner ignores store: false 3 a test fails for each
a blocked or an interrupted turn is left to OGX 2 a test fails for each

28 of the 30 fail a test. The two that do not:

  • The check after the write of a blocked /v1/responses stream. No test takes that write out, so nothing pins the check. Removing the write itself there fails test_blocked_turn_is_stored_once[compacted-streaming]. The approved commit has the same three lines and the same tests around them.
  • The interrupt guard at the end of the stream. Nothing changes without it: when the guard is set there, the owner has settled the turn, and a later report does nothing.
  1. Check on a running service that every turn is stored once. This was run again on the rebased commit on 2026-10-07, in library mode; the output is from that run. It sent the five queries the e2e compaction scenario had before LCORE-4219: keep the buffered turns after a compaction #2791, which changed the second and the fifth. Start the service in library mode with compaction enabled (context_windows: openai/gpt-4o-mini: 2000, threshold_ratio: 0.1, token_floor: 100, buffer_turns: 1), send the five queries of the e2e compaction scenario through each endpoint, and read the stored items from the OGX database:
/v1/query
  query 1: context_status=full response='OK'
  query 2: context_status=full response='OK'
  query 3: context_status=summarized response='aurora-prod-7, blue-lagoon'
  query 4: context_status=summarized response='OK'
  query 5: context_status=summarized response='north-quarry, green-harbor'
  stored: 13 items = 5 queries + 5 answers + 3 markers
  every query stored once: True; query/answer alternate: True
/v1/streaming_query
  ...
  stored: 13 items = 5 queries + 5 answers + 3 markers
  every query stored once: True; query/answer alternate: True
/v1/responses
  ...
  stored: 13 items = 5 queries + 5 answers + 3 markers
  every query stored once: True; query/answer alternate: True
/v1/responses (stream)
  ...
  stored: 13 items = 5 queries + 5 answers + 3 markers
  every query stored once: True; query/answer alternate: True
RESULT: every turn stored exactly once on every endpoint

A2A was not run live; its integration tests cover it.

  1. Run the full suites:
uv run pytest tests/unit -q                                                          # 3735 passed, 1 skipped
uv run pytest tests/integration --ignore=tests/integration/container_lifecycle -q   # 370 passed (the ignored module builds and starts the OGX container; not touched by this PR)
  1. Linters:
uv run make black ruff docstyle pylint pyright   # all pass
uv run make check-types                          # mypy fails only where main does: one error in src/models/config.py, the rest in tests, none in a line this PR changes

Summary by CodeRabbit

  • Bug Fixes
    • Improved conversation history handling across query, response, streaming, and A2A requests. Completed, blocked, and interrupted turns are stored consistently without duplicate entries.
    • Compacted requests retain the original user input in conversation history.
    • Persistence failures and turns that are not stored are handled consistently across endpoint types.

@coderabbitai

coderabbitai Bot commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

Warning

Review limit reached

You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository.

Next included review available in 59 minutes.

Check out review usage here.

View limit details

Limit details: You’ve used the included review currently available.

Learn how review limits work.

Review configuration:

⚙️ Run configuration
  • Configuration used: Repository: lightspeed-core/lightspeed-stack/.coderabbit.yaml
  • Review profile: ASSERTIVE
  • Plan: Advanced
  • Run ID: b1d6c419-b3ff-4969-a3e6-ad72120b744b
📥 Commits

Reviewing files that changed from the base of the PR and between eb13198 and 3f919be.

📒 Files selected for processing (23)
  • docs/design/conversation-compaction/conversation-compaction.md
  • docs/devel_doc/ARCHITECTURE.md
  • src/app/endpoints/a2a.py
  • src/app/endpoints/query.py
  • src/app/endpoints/responses.py
  • src/app/endpoints/streaming_query.py
  • src/utils/agents/query.py
  • src/utils/agents/streaming.py
  • src/utils/conversation_compaction.py
  • src/utils/pending_turn.py
  • src/utils/stream_interrupts.py
  • tests/integration/endpoints/_compaction_helpers.py
  • tests/integration/endpoints/test_compaction_a2a.py
  • tests/integration/endpoints/test_turn_persistence.py
  • tests/unit/app/endpoints/test_a2a.py
  • tests/unit/app/endpoints/test_query.py
  • tests/unit/app/endpoints/test_query_otel.py
  • tests/unit/app/endpoints/test_responses.py
  • tests/unit/utils/agents/test_query.py
  • tests/unit/utils/agents/test_streaming.py
  • tests/unit/utils/test_conversation_compaction.py
  • tests/unit/utils/test_pending_turn.py
  • tests/unit/utils/test_stream_interrupts.py

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration
  • Configuration used: Repository: lightspeed-core/lightspeed-stack/.coderabbit.yaml
  • Review profile: ASSERTIVE
  • Plan: Advanced
  • Run ID: b7aee751-73ce-4f10-a636-7764628a25af
📥 Commits

Reviewing files that changed from the base of the PR and between 5196e2a and eb13198.

📒 Files selected for processing (14)
  • src/app/endpoints/a2a.py
  • src/app/endpoints/query.py
  • src/app/endpoints/responses.py
  • src/app/endpoints/streaming_query.py
  • src/utils/agents/query.py
  • src/utils/agents/streaming.py
  • src/utils/conversation_compaction.py
  • src/utils/stream_interrupts.py
  • tests/integration/endpoints/test_compaction_a2a.py
  • tests/unit/app/endpoints/test_a2a.py
  • tests/unit/app/endpoints/test_query.py
  • tests/unit/app/endpoints/test_responses.py
  • tests/unit/utils/agents/test_query.py
  • tests/unit/utils/agents/test_streaming.py
💤 Files with no reviewable changes (1)
  • src/utils/conversation_compaction.py

Included review availability: This review used your included allowance. Your plan provides up to 1 included review per hour; 0 remain after this review.

📜 Recent review details
⏰ Context from checks skipped due to timeout. (26)
  • GitHub Check: E2E: library / ci / skills
  • GitHub Check: E2E: library / ci / default
  • GitHub Check: E2E: server / ci / tls
  • GitHub Check: E2E: library / ci / rbac
  • GitHub Check: E2E: library / ci / other
  • GitHub Check: E2E: library / ci / authorized
  • GitHub Check: E2E: server / ci / mcp
  • GitHub Check: E2E: server / ci / rbac
  • GitHub Check: E2E: server / ci / default
  • GitHub Check: E2E: library / ci / mcp
  • GitHub Check: E2E: library / ci / shields
  • GitHub Check: E2E: server / ci / other
  • GitHub Check: E2E: server / ci / skills
  • GitHub Check: E2E: server / ci / shields
  • GitHub Check: E2E: server / ci / authorized
  • GitHub Check: spectral
  • GitHub Check: build-pr
  • GitHub Check: unit_tests (3.13)
  • GitHub Check: integration_tests (3.12)
  • GitHub Check: Pylinter
  • GitHub Check: unit_tests (3.12)
  • GitHub Check: integration_tests (3.13)
  • GitHub Check: Red Hat Konflux / lightspeed-stack-0-8-e2e-tests / lightspeed-stack-0-8
  • GitHub Check: Red Hat Konflux / rag-content-0-8-e2e-tests / lightspeed-stack-0-8
  • GitHub Check: Red Hat Konflux / lightspeed-core-0-8-enterprise-contract / lightspeed-stack-0-8
  • GitHub Check: Konflux kflux-prd-rh02 / lightspeed-stack-0-8-on-pull-request
🔇 Additional comments (5)
src/utils/stream_interrupts.py (1)

9-9: LGTM!

Also applies to: 20-20, 366-366

src/utils/agents/query.py (1)

49-49: LGTM!

src/utils/agents/streaming.py (1)

75-75: LGTM!

tests/unit/utils/agents/test_query.py (1)

37-37: LGTM!

Also applies to: 355-366, 424-424, 427-427, 430-430, 436-464, 668-668, 674-674, 699-699

tests/unit/utils/agents/test_streaming.py (1)

70-70: LGTM!

Also applies to: 493-505, 1615-1615, 1619-1620, 1640-1640, 1643-1645, 1665-1665, 1695-1695, 1697-1700, 1703-1710, 1731-1731, 1734-1734, 1739-1743, 1745-1797, 1804-1804, 1808-1808, 1832-1832, 1835-1835


Walkthrough

This change adds PendingTurn to coordinate conversation-turn ownership, storage, and settlement. Query, responses, streaming query, A2A, agent, and interruption paths use it. Tests cover persistence outcomes, write failures, and skipped storage.

Changes

Turn persistence ownership

Layer / File(s) Summary
PendingTurn contract and settlement
src/utils/pending_turn.py, src/utils/conversation_compaction.py, tests/unit/utils/test_pending_turn.py, tests/unit/utils/test_conversation_compaction.py, docs/design/conversation-compaction/conversation-compaction.md, docs/devel_doc/ARCHITECTURE.md
PendingTurn selects request input and storage ownership, records the first outcome, prevents duplicate writes, and checks settlement. Compacted requests require original input. The store_compacted_turn helper is removed, and the design and architecture documents describe the ownership model.
Endpoint turn wiring
src/app/endpoints/*, tests/unit/app/endpoints/*
Query, responses, streaming query, and A2A handlers create or receive a pending turn and use it for persistence and settlement. Endpoint tests update fixtures and writer patches.
Agent and interruption persistence
src/utils/agents/*, src/utils/stream_interrupts.py, tests/unit/utils/agents/*, tests/unit/utils/test_stream_interrupts.py
Agent response and interruption paths pass a shared pending turn instead of original input. Tests cover compacted input storage, interrupted outcomes, and duplicate prevention.
Endpoint persistence validation
tests/integration/endpoints/test_turn_persistence.py, tests/integration/endpoints/_compaction_helpers.py, tests/integration/endpoints/test_compaction_a2a.py
Integration tests cover persistence outcomes across endpoints, including failed writes and skipped storage. A2A compaction tests use shared request, agent, and card helpers.

Priority: ➖ Normal

Estimated code review effort: 4 (Complex) | ~45 minutes

Change: Refactor

Sequence Diagram(s)

sequenceDiagram
  participant Endpoint
  participant Agent
  participant PendingTurn
  participant ConversationStore
  Endpoint->>PendingTurn: create turn from request parameters and original input
  Endpoint->>Agent: run with shared turn
  Agent->>PendingTurn: report completed or interrupted outcome
  PendingTurn->>ConversationStore: append owned turn once
  Endpoint->>PendingTurn: check settlement
Loading

Suggested reviewers: asimurka, tisnik

Merge Risk: ⚪ Minimal · up to eb131

This change centralizes how conversation turns are stored. No actionable merge-blocking risk was identified in the reviewed material.

Security Architecture Review

Security architecture risk: 🟡 Moderate · up to 2af20

Conversation history now depends on one shared persistence flow across several request types. The reviewed paths retain their existing conversation-access checks, and no new security exposure was verified. Failed writes and interrupted execution can still leave a turn out of later context, so this is a meaningful design change with residual uncertainty.

Retained concerns
No architecture-level concerns identified.

Security review details

Security Blast Radius

  • inferred — Requests able to use the affected query, Responses, streaming, or A2A paths can affect persistence of their resolved conversation turns. The inspected owner does not itself grant access to another conversation or introduce a new external sink.

Trust Boundaries and Controls

  • observed — On the inspected Responses path, conversation access is resolved before moderation and before PendingTurn receives prepared parameters; previous_response_id changes storage ownership rather than authorizing a conversation inside the owner.

Resilience and Maintainability Implications

  • inferred — First settlement contains duplicate writes from competing lifecycle callbacks on one owner, but it is not a recovery mechanism: a failed or cancelled append can leave later conversation context without the turn. The inspected comparison does not establish this as a new PR-introduced exposure.

Hardening Proposals

  • proposed — If durable conversation history is required for a security control, distinguish an attempted append from a confirmed write and define recovery for failed or cancelled writes without allowing duplicate turns.
🚥 Pre-merge checks | ✅ 7
✅ Passed checks (7 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly summarizes the main change: centralizing ownership of turns stored by the service through a single persistence owner.
Docstring Coverage ✅ Passed Docstring coverage is 94.48% which is sufficient. The required threshold is 80.00%. Docstring coverage is scoped to functions touched by this diff. Analyzed 181 functions across 21 files.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Performance And Algorithmic Complexity ✅ Passed No meaningful performance regression is introduced. PendingTurn performs at most one conversation write per request: _settle prevents duplicate writes before the single `append_turn_items_to_conve…
Security And Secret Handling ✅ Passed No security violation was introduced. The changed API routes retain their existing authentication dependencies and @authorize checks. PendingTurn only forwards validated conversation inputs and mo…
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create a new PR
✨ Simplify code
  • Create a new PR

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@tisnik tisnik left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM but there are merge conflicts

@max-svistunov
max-svistunov force-pushed the lcore-3908-compacted-turn-persistence-owner branch from 2af20de to 5196e2a Compare October 7, 2026 20:10
max-svistunov added a commit to max-svistunov/lightspeed-stack that referenced this pull request Oct 7, 2026
… approval

The parent of this commit is PR lightspeed-core#2796 as it was approved on 2026-09-30,
re-applied onto current main with only the conflicts resolved. The tree
of this commit is the branch of the PR as it is now. So the diff of
this one commit is everything that changed since the approval:

- main removed the shield moderation that ran before the agent on
  /v1/query and /v1/streaming_query (lightspeed-core#2780), so what the PR had written
  against it is gone: PendingTurn.drop(), the turn parameter of
  retrieve_agent_response_generator, the turn the streaming handler
  created for it, the tests that pinned that path (five integration
  tests, one unit test) and three rows of the design-doc table;
- the A2A span reads the compacted flag from the request parameters
  (main's span code read a variable the PR had removed);
- the scan of src for conversation writers also lists
  replace_last_assistant_message, which main added;
- two docstrings and one doc sentence follow the above.
@max-svistunov
max-svistunov force-pushed the lcore-3908-compacted-turn-persistence-owner branch from 5196e2a to eb13198 Compare October 7, 2026 22:50
OGX appends a turn to the conversation only when it is handed the
conversation parameter and runs the inference. In compacted mode that
parameter is dropped, so lightspeed-stack appends the turn itself. Every
endpoint made that write on its own, and nothing enforced it: two
unrelated cleanups removed the calls, and for six weeks conversations
stopped growing after a compaction without anything failing
(LCORE-3883).

The write now has one owner, PendingTurn in utils/pending_turn.py. An
endpoint creates it from the parameters the request is sent with and
the original input of the compaction result, and reports how the turn
ended: store_completed, store_blocked or store_interrupted. The owner
decides whether the turn is ours, which input is stored, and it stores
a turn once: the first report settles the turn, every later one does
nothing.

It covers all seven places that stored a turn, not only the compacted
ones, because they are the same duty: a turn OGX does not store. That
is a conversation in compacted mode, a request a shield blocked on
/v1/responses, a stream the client interrupted, a continuation from
previous_response_id. store_compacted_turn is removed; no endpoint
appends a turn except through the owner. The shield capabilities still
store the turn they rejected themselves, from inside the agent run, and
only when the model was handed the conversation, so never in compacted
mode. /v1/query and /v1/streaming_query have no blocked turn of their
own to store: the shield moderation that ran before the agent on those
two endpoints is gone from main (lightspeed-core#2780).

A compacted request that loses its turn through a change in the code
fails:
- creating the owner for compacted parameters without the original
  input raises ValueError, which is how a caller that lost the
  hand-over looks;
- ensure_settled raises TurnNotStoredError when a compacted request
  reaches the end of its handler and nobody tried to store its turn.
  /v1/query runs in a scope that makes the check when it is left; the
  streaming paths, /v1/responses and A2A make it after their write, in
  the function that calls the write, not in the one that performs it.
  The error is not a RuntimeError, because the endpoints report that
  one as an inference failure.
The check does not cover a stream the client stops reading, a
/v1/responses stream without a final response, and a failed write that
is logged. They store nothing, as before.

Every condition of the old call sites is kept: where the write
happens, what is stored for a failed or incomplete turn, what a failed
write does to the request (it fails /v1/query and /v1/responses, it is
logged on /v1/streaming_query, A2A and the interrupt path), and the
order relative to quota consumption. The guard that lets only one of
stream end, cancellation handler and interrupt callback finish a turn
is unchanged. What is stored, and when, does not change.

Tests:
- integration, tests/integration/endpoints/test_turn_persistence.py:
  for /v1/query, /v1/streaming_query, /v1/responses (blocking and
  streaming) and A2A, every way a turn can end (completed, blocked on
  /v1/responses, interrupted, failed, cut short by the model, abandoned
  by the client, continued from a previous response, store off),
  compacted and not; what a failed write does to the request on each
  endpoint; and what happens when the step that stores the turn is
  taken out of an endpoint. Each test runs the real handler and
  compares the conversation item by item afterwards, so a missing
  write, a second write and a write of the wrong input all fail it. 39
  of the 44 tests pass on main unchanged. The five that fail there take
  the write out of an endpoint, which main does not notice.
- unit: the owner, the hand-over from the compaction seam
  (omit_conversation and the original input are set together), and a
  scan of src that fails when a module other than the known writers
  mentions one of the functions that write to a conversation.
- existing unit tests follow the new parameter. Those that asserted on
  the arguments of a patched helper in the agent and interrupt paths
  now assert on what is stored. Mocked request parameters got the
  fields the owner reads. The A2A helpers of the compaction tests moved
  to _compaction_helpers.py, where both test modules import them.

The tests were checked by mutation: with the write removed from any
of the four endpoints, a check removed, the write doubled, the settled
state ignored, the interrupt guard ignored, the explicit rewrite stored
in place of the input, a turn OGX stores stored by us as well, or the
store flag ignored, tests fail. One check is not pinned by a test: the
one after the write of a blocked /v1/responses stream.

Docs: the design document describes the owner and has a table of what
is stored per endpoint and per way a turn can end.
@max-svistunov
max-svistunov force-pushed the lcore-3908-compacted-turn-persistence-owner branch from eb13198 to 3f919be Compare October 8, 2026 16:00

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants