Repository navigation
LCORE-3908: give the turns we store ourselves one owner - #2796
max-svistunov wants to merge 1 commit into
Conversation
|
Warning Review limit reachedYou'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. View limit detailsLimit details: You’ve used the included review currently available. Review configuration: ⚙️ Run configuration
📒 Files selected for processing (23)
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configuration
📒 Files selected for processing (14)
💤 Files with no reviewable changes (1)
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)
🔇 Additional comments (5)
WalkthroughThis change adds ChangesTurn persistence ownership
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
Suggested reviewers: Merge Risk: ⚪ Minimal · up to This change centralizes how conversation turns are stored. No actionable merge-blocking risk was identified in the reviewed material. Security Architecture ReviewSecurity architecture risk: 🟡 Moderate · up to 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 Security review detailsSecurity Blast Radius
Trust Boundaries and Controls
Resilience and Maintainability Implications
Hardening Proposals
🚥 Pre-merge checks | ✅ 7✅ Passed checks (7 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
✨ Simplify code
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. Comment |
tisnik
left a comment
There was a problem hiding this comment.
LGTM but there are merge conflicts
2af20de to
5196e2a
Compare
… 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.
5196e2a to
eb13198
Compare
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.
eb13198 to
3f919be
Compare
Description
LCORE-3908. OGX appends a turn to the conversation only when it is handed the
conversationparameter 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 wheremainchanged 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.mainremoved the shield moderation that ran before the agent on/v1/queryand/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 (mainremoved it).retrieve_agent_responseis called withoutmoderation_result, as onmain.retrieve_agent_response_generatorloses theturnparameter and is as onmain.src/app/endpoints/streaming_query.pyno longer creates aPendingTurnfor the not-compacted path: it was there to share one owner between the refusal written before the stream and the stream itself.generate_agent_responsecreates the owner from the parameters, as it already did when none was passed.PendingTurn.drop()is removed with its only caller, the/v1/querybranch, and so is the logger it used. The docstrings, a comment insrc/app/endpoints/query.pyand the message ofTurnNotStoredErrorno 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_laterand the dropped case oftest_settled_turn_passes_the_checkare deleted, and three docstrings are reworded.tests/unit/utils/agents/: the edits to the four tests of the blocked branches are gone, becausemaindeleted those tests.tests/unit/app/endpoints/test_streaming_query.pyis no longer touched. Its only changes were mock fields for the owner the handler created./v1/queryand/v1/streaming_queryand the row "blocked by a shield, then interrupted", with the paragraph about that row. The text around the table no longer mentionsdrop, 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.mdis no longer touched: its one sentence described the removed path.mainrecords on the A2A span whether the turn was compacted, and reads it fromcompaction.compacted(LCORE-3755: Emit additional Otel attributes for evaluation #2818). This PR moves the compaction call into a helper, sosrc/app/endpoints/a2a.pypassesresponses_params.omit_conversationthere.omit_conversationis set in one place, together withcompacted=True.mainaddedreplace_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 ofsrcintests/unit/utils/test_pending_turn.pylists it, with the three modules that mention it.tests/unit/app/endpoints/test_a2a.py: one two-line hunk is gone, becausemainhas the same two lines now (LCORE-3755: Emit additional Otel attributes for evaluation #2818).src/app/endpoints/responses.pyis one of the 11:mainchanged that file too (#2780, #2781, #2818), and the same lines sit in the new code.The owner. The write now has one owner,
PendingTurninsrc/utils/pending_turn.py. An endpoint creates it from the parameters the request is sent with and theoriginal_inputof theCompactionResult, and reports how the turn ended:store_completed,store_blockedorstore_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 fromprevious_response_id.store_compacted_turnis removed. No endpoint appends a turn except through the owner.A compacted request that loses its turn through a change in the code fails.
ValueError. That is how a caller that lost the hand-over looks.ensure_settled()raisesTurnNotStoredErrorwhen a compacted request reaches the end of its handler and nobody tried to store its turn./v1/queryruns in thepending_turn()scope, which makes the check when it is left. The streaming paths,/v1/responsesand A2A make it after their write, in the function that calls the write and not in the one that performs it.TurnNotStoredErroris not aRuntimeError, 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.
ValueErrorin compacted mode; outside compacted mode the helper creates an owner of its own, which is what each call site did before.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/queryand/v1/streaming_queryhave no blocked turn of their own to store: the shield moderation that ran before the agent on those two endpoints is gone frommain(LCORE-4338: Remove obsolete pre-agent shield moderation from query and streaming. #2780)./v1/responsesstream without a final response, a failed write that is logged. Tests record all three.src/utils/pending_turn.pyPendingTurn,pending_turn(),TurnNotStoredErrorsrc/utils/agents/query.py,src/app/endpoints/query.pyoriginal_input; the scopesrc/utils/agents/streaming.py,src/utils/stream_interrupts.py,src/app/endpoints/streaming_query.pyoriginal_input; the checksrc/app/endpoints/responses.pysrc/app/endpoints/a2a.pysrc/utils/conversation_compaction.pystore_compacted_turnremovedtests/integration/endpoints/test_turn_persistence.pytests/unit/utils/test_pending_turn.pytests/unit/...tests/integration/endpoints/_compaction_helpers.py,test_compaction_a2a.pydocs/design/conversation-compaction/conversation-compaction.md,docs/devel_doc/ARCHITECTURE.mdRelation 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.pyis at 890 lines with this change, and pylint stops at 1000.Type of change
pyproject.toml+uv.lock]requirements.*.txtfor Konflux]Tools used to create PR
Identify any AI code assistants used in this PR (for transparency and review context)
Related Tickets & Documents
Checklist before requesting a review
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 forreplace_last_assistant_messagein the scan ofsrc, andtests/unit/app/endpoints/test_streaming_query.py, which is back as onmain. Step 3 was run on the rebased commit as it is now, on 2026-10-07, in library mode.The same file against the
src/ofmainshows that behaviour is kept (with the test helpers it imports, andsrc/utils/pending_turn.pyadded so that it can be collected). 39 pass. The five that fail take the write out of an endpoint, whichmaindoes not notice: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.pyandtests/unit/app/endpoints/test_responses.py(298 tests).src//v1/query,/v1/streaming_query, the interrupt path, A2A,/v1/responses(completed and blocked, blocking and streaming)/v1/streaming_query, A2A,/v1/responses(blocking, streaming), theValueErrorfor a missing original input/v1/responsesstreamstore: false28 of the 30 fail a test. The two that do not:
/v1/responsesstream. No test takes that write out, so nothing pins the check. Removing the write itself there failstest_blocked_turn_is_stored_once[compacted-streaming]. The approved commit has the same three lines and the same tests around them.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:A2A was not run live; its integration tests cover it.
Summary by CodeRabbit