From 7be87bfc9196358c099e90025b19fbcb7cb2e33a Mon Sep 17 00:00:00 2001 From: kl3inIT Date: Thu, 6 Aug 2026 15:44:09 +0700 Subject: [PATCH] feat(assistant): give a conversation turn an explicit identity A turn writes its question in beginTurn and its answer in completeTurn, in separate transactions. Two turns of one conversation can therefore persist as U1, U2, A2, A1, and no ordering heuristic over sequence_id pairs them. Nothing reads the pairing today, but the transcript context reader that replaces Spring AI's chat memory must, and it cannot recover from sequence order what the writers already knew. beginTurn now allocates a turn id and returns it with the conversation id; completeTurn takes that reference and writes the answer under the same identity. A partial unique index over (turn_id, role) holds one question and one answer per turn. Rows written before this migration keep a null turn_id: they stay visible in the transcript and are exempt from the index, because their pairing cannot be recovered after the fact. Two turns completing at the same instant still raise an optimistic-locking failure through the unlocked conversation touch in completeTurn. That is pre-existing and unrelated to turn identity; it is recorded as a gap in the Assistant test matrix rather than fixed here. Co-Authored-By: Claude Opus 5 (1M context) --- .../api/assistant/AssistantController.java | 10 +- ...erFeedbackConcurrencyIntegrationTests.java | 5 +- .../AssistantControllerStreamingTests.java | 9 +- ...lSelectionConcurrencyIntegrationTests.java | 3 +- ...AssistantTurnIdentityIntegrationTests.java | 210 ++++++++++++++++++ .../AssistantConversationMessage.java | 10 + .../AssistantConversationService.java | 17 +- .../core/assistant/AssistantTurnRef.java | 21 ++ .../V26__assistant_message_turn_identity.sql | 20 ++ .../AssistantConversationServiceTests.java | 55 ++++- .../plan.md | 4 +- docs/specs/domains/assistant-and-mcp.md | 11 +- docs/tests/domains/assistant-and-mcp.md | 4 +- 13 files changed, 356 insertions(+), 23 deletions(-) create mode 100644 apps/api/src/test/java/com/orgmemory/api/assistant/AssistantTurnIdentityIntegrationTests.java create mode 100644 core/src/main/java/com/orgmemory/core/assistant/AssistantTurnRef.java create mode 100644 core/src/main/resources/db/migration/V26__assistant_message_turn_identity.sql diff --git a/apps/api/src/main/java/com/orgmemory/api/assistant/AssistantController.java b/apps/api/src/main/java/com/orgmemory/api/assistant/AssistantController.java index be843f49..d5dcc4ae 100644 --- a/apps/api/src/main/java/com/orgmemory/api/assistant/AssistantController.java +++ b/apps/api/src/main/java/com/orgmemory/api/assistant/AssistantController.java @@ -10,6 +10,7 @@ import com.orgmemory.core.assistant.AssistantConversationSummary; import com.orgmemory.core.assistant.AssistantService; import com.orgmemory.core.assistant.AssistantTurn; +import com.orgmemory.core.assistant.AssistantTurnRef; import com.orgmemory.core.ai.AssistantModelAuthorityService; import com.orgmemory.core.ai.AssistantModelChoice; import com.orgmemory.core.ai.AssistantModelRouteAuthority; @@ -111,11 +112,12 @@ ResponseEntity>> chat( request.modelActivationId()); AssistantModelSelectionRef modelSelection = modelAuthority.selectionRef( routeAuthority); - UUID conversationId = conversations.beginTurn( + AssistantTurnRef turnRef = conversations.beginTurn( actor, request.conversationId(), request.message(), modelSelection); + UUID conversationId = turnRef.conversationId(); UUID assistantMessageId = UUID.randomUUID(); Flux parts = Flux.defer(() -> { long turnStartedAtNanos = System.nanoTime(); @@ -136,7 +138,7 @@ ResponseEntity>> chat( turnStartedAtNanos)) .flatMapMany(turn -> completedTurnParts( actor, - conversationId, + turnRef, assistantMessageId, turn))); }); @@ -418,7 +420,7 @@ Flux parts(AssistantTurn turn) { private Flux completedTurnParts( CurrentActor actor, - UUID conversationId, + AssistantTurnRef turnRef, UUID assistantMessageId, AssistantTurn turn) { StringBuilder completedAnswer = new StringBuilder(); @@ -430,7 +432,7 @@ private Flux completedTurnParts( }) .doOnComplete(() -> conversations.completeTurn( actor, - conversationId, + turnRef, assistantMessageId, completedAnswer.toString(), turn.citations())); diff --git a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantAnswerFeedbackConcurrencyIntegrationTests.java b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantAnswerFeedbackConcurrencyIntegrationTests.java index cc3c8008..274b50cf 100644 --- a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantAnswerFeedbackConcurrencyIntegrationTests.java +++ b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantAnswerFeedbackConcurrencyIntegrationTests.java @@ -4,6 +4,7 @@ import com.orgmemory.core.assistant.AssistantAnswerSentiment; import com.orgmemory.core.assistant.AssistantConversationService; +import com.orgmemory.core.assistant.AssistantTurnRef; import com.orgmemory.core.organization.CurrentActor; import com.orgmemory.core.organization.Clearance; import java.util.ArrayList; @@ -116,9 +117,9 @@ INSERT INTO app_users ( "Feedback actor", actorId + "@example.test", Clearance.STANDARD); - UUID conversationId = conversations.beginTurn(actor, null, "What is the policy?"); + AssistantTurnRef turn = conversations.beginTurn(actor, null, "What is the policy?"); UUID answerId = UUID.randomUUID(); - conversations.completeTurn(actor, conversationId, answerId, "The policy is available."); + conversations.completeTurn(actor, turn, answerId, "The policy is available."); return new Scenario(actor, answerId); } diff --git a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantControllerStreamingTests.java b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantControllerStreamingTests.java index 03456817..0624f912 100644 --- a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantControllerStreamingTests.java +++ b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantControllerStreamingTests.java @@ -23,6 +23,7 @@ import com.orgmemory.core.assistant.AssistantConversationService; import com.orgmemory.core.assistant.AssistantService; import com.orgmemory.core.assistant.AssistantTurn; +import com.orgmemory.core.assistant.AssistantTurnRef; import com.orgmemory.core.knowledge.search.RetrievedKnowledgeEvidence; import com.orgmemory.core.knowledge.retrieval.CitationEvidenceService; import com.orgmemory.core.knowledge.retrieval.CitationEvidenceReference; @@ -226,8 +227,9 @@ void usesOneServerOwnedIdentityForTheStreamAndPersistedAnswer() { "laura@example.test"); UUID conversationId = UUID.randomUUID(); when(actors.current(authentication)).thenReturn(actor); + AssistantTurnRef turnRef = new AssistantTurnRef(conversationId, UUID.randomUUID()); when(conversations.beginTurn(actor, null, "Question", null)) - .thenReturn(conversationId); + .thenReturn(turnRef); when(assistant.startTurn( eq(actor), eq("Question"), @@ -262,7 +264,7 @@ void usesOneServerOwnedIdentityForTheStreamAndPersistedAnswer() { ArgumentCaptor messageId = ArgumentCaptor.forClass(UUID.class); verify(conversations).completeTurn( eq(actor), - eq(conversationId), + eq(turnRef), messageId.capture(), eq("Answer"), eq(List.of())); @@ -292,8 +294,9 @@ void emitsStreamStartAndRetrievalActivityWhileRetrievalIsStillBlocked() CountDownLatch releaseRetrieval = new CountDownLatch(1); AtomicBoolean retrievalCompleted = new AtomicBoolean(); when(actors.current(authentication)).thenReturn(actor); + AssistantTurnRef turnRef = new AssistantTurnRef(conversationId, UUID.randomUUID()); when(conversations.beginTurn(actor, null, "Question", null)) - .thenReturn(conversationId); + .thenReturn(turnRef); when(assistant.startTurn( eq(actor), eq("Question"), diff --git a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantModelSelectionConcurrencyIntegrationTests.java b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantModelSelectionConcurrencyIntegrationTests.java index ec7870b9..3080dcb3 100644 --- a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantModelSelectionConcurrencyIntegrationTests.java +++ b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantModelSelectionConcurrencyIntegrationTests.java @@ -181,7 +181,8 @@ INSERT INTO app_users ( "Model actor", actorId + "@example.test", Clearance.STANDARD); - UUID conversationId = conversations.beginTurn(actor, null, "Initial turn"); + UUID conversationId = + conversations.beginTurn(actor, null, "Initial turn").conversationId(); return new Scenario(actor, profile.id(), activation.id(), conversationId); } diff --git a/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantTurnIdentityIntegrationTests.java b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantTurnIdentityIntegrationTests.java new file mode 100644 index 00000000..c623edc5 --- /dev/null +++ b/apps/api/src/test/java/com/orgmemory/api/assistant/AssistantTurnIdentityIntegrationTests.java @@ -0,0 +1,210 @@ +package com.orgmemory.api.assistant; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; + +import com.orgmemory.core.assistant.AssistantConversationService; +import com.orgmemory.core.assistant.AssistantTurnRef; +import com.orgmemory.core.organization.Clearance; +import com.orgmemory.core.organization.CurrentActor; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.testcontainers.service.connection.ServiceConnection; +import org.springframework.dao.DuplicateKeyException; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.annotation.DirtiesContext; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.postgresql.PostgreSQLContainer; + +@SpringBootTest +@Testcontainers +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) +class AssistantTurnIdentityIntegrationTests { + + @Container + @ServiceConnection + static PostgreSQLContainer postgres = new PostgreSQLContainer("pgvector/pgvector:pg18"); + + @Autowired + AssistantConversationService conversations; + + @Autowired + JdbcTemplate jdbc; + + @Test + void pairsTheQuestionAndAnswerOfOneTurnUnderOneIdentity() { + CurrentActor actor = seededActor(); + + AssistantTurnRef turn = conversations.beginTurn(actor, null, "How long is probation?"); + conversations.completeTurn(actor, turn, UUID.randomUUID(), "Sixty days."); + + assertEquals( + List.of("USER", "ASSISTANT"), + jdbc.queryForList( + """ + SELECT role FROM assistant_conversation_messages + WHERE turn_id = ? ORDER BY sequence_id + """, + String.class, + turn.turnId())); + } + + @Test + void pairsOverlappingTurnsThatPersistOutOfOrder() { + CurrentActor actor = seededActor(); + UUID conversationId = conversations.beginTurn(actor, null, "Opening turn") + .conversationId(); + + // Two turns open before either answers, and the second answers first: + // U1, U2, A2, A1. No ordering heuristic over sequence_id pairs this. + AssistantTurnRef first = conversations.beginTurn(actor, conversationId, "First question"); + AssistantTurnRef second = conversations.beginTurn(actor, conversationId, "Second question"); + conversations.completeTurn(actor, second, UUID.randomUUID(), "Second answer"); + conversations.completeTurn(actor, first, UUID.randomUUID(), "First answer"); + + assertEquals( + List.of("First question", "First answer"), + contentOfTurn(first.turnId())); + assertEquals( + List.of("Second question", "Second answer"), + contentOfTurn(second.turnId())); + } + + @Test + void givesEveryConcurrentQuestionInOneConversationItsOwnIdentity() throws Exception { + CurrentActor actor = seededActor(); + UUID conversationId = conversations.beginTurn(actor, null, "Opening turn") + .conversationId(); + int turnCount = 12; + + CountDownLatch ready = new CountDownLatch(turnCount); + CountDownLatch start = new CountDownLatch(1); + List> started = new ArrayList<>(); + try (var executor = Executors.newFixedThreadPool(turnCount)) { + for (int index = 0; index < turnCount; index++) { + int turn = index; + started.add(executor.submit(() -> { + ready.countDown(); + start.await(10, TimeUnit.SECONDS); + return conversations.beginTurn( + actor, conversationId, "Concurrent question " + turn); + })); + } + ready.await(10, TimeUnit.SECONDS); + start.countDown(); + for (Future attempt : started) { + assertNotNull(attempt.get(30, TimeUnit.SECONDS)); + } + } + + assertEquals( + turnCount + 1, + jdbc.queryForObject( + """ + SELECT count(DISTINCT turn_id) FROM assistant_conversation_messages + WHERE conversation_id = ? + """, + Integer.class, + conversationId)); + } + + @Test + void refusesASecondMessageOfTheSameRoleInOneTurn() { + CurrentActor actor = seededActor(); + AssistantTurnRef turn = conversations.beginTurn(actor, null, "How long is probation?"); + + assertThrows( + DuplicateKeyException.class, + () -> insertMessage(actor, turn.conversationId(), turn.turnId(), "USER", "Again")); + } + + @Test + void leavesRowsWrittenBeforeTurnIdentityExistedUnconstrained() { + CurrentActor actor = seededActor(); + UUID conversationId = conversations.beginTurn(actor, null, "Opening turn") + .conversationId(); + + insertMessage(actor, conversationId, null, "USER", "Legacy question"); + insertMessage(actor, conversationId, null, "USER", "Another legacy question"); + + // Partial uniqueness: legacy rows carry no pairing to protect, so they + // must not collide with one another and must stay in the transcript. + assertEquals( + 2, + jdbc.queryForObject( + """ + SELECT count(*) FROM assistant_conversation_messages + WHERE conversation_id = ? AND turn_id IS NULL + """, + Integer.class, + conversationId)); + } + + private List contentOfTurn(UUID turnId) { + return jdbc.queryForList( + """ + SELECT content FROM assistant_conversation_messages + WHERE turn_id = ? ORDER BY sequence_id + """, + String.class, + turnId); + } + + private void insertMessage( + CurrentActor actor, UUID conversationId, UUID turnId, String role, String content) { + jdbc.update( + """ + INSERT INTO assistant_conversation_messages ( + id, conversation_id, turn_id, organization_id, actor_user_id, + role, content, occurred_at, created_at, updated_at, version) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, now(), now(), 0) + """, + UUID.randomUUID(), + conversationId, + turnId, + actor.organizationId(), + actor.userId(), + role, + content, + java.sql.Timestamp.from(Instant.now())); + } + + private CurrentActor seededActor() { + UUID organizationId = UUID.randomUUID(); + UUID actorId = UUID.randomUUID(); + jdbc.update( + """ + INSERT INTO organizations (id, name, created_at, updated_at, version) + VALUES (?, 'Turn identity', now(), now(), 0) + """, + organizationId); + jdbc.update( + """ + INSERT INTO app_users ( + id, organization_id, name, email, clearance, active, + created_at, updated_at, version) + VALUES (?, ?, 'Turn identity actor', ?, 'STANDARD', true, now(), now(), 0) + """, + actorId, + organizationId, + actorId + "@example.test"); + return new CurrentActor( + actorId, + organizationId, + null, + "Turn identity actor", + actorId + "@example.test", + Clearance.STANDARD); + } +} diff --git a/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationMessage.java b/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationMessage.java index 412f6c5d..49126c97 100644 --- a/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationMessage.java +++ b/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationMessage.java @@ -17,6 +17,10 @@ class AssistantConversationMessage extends BaseEntity { @Column(name = "conversation_id", nullable = false, updatable = false) private UUID conversationId; + /** Null only on rows written before turn identity existed. */ + @Column(name = "turn_id", updatable = false) + private UUID turnId; + @Column(name = "organization_id", nullable = false, updatable = false) private UUID organizationId; @@ -42,6 +46,7 @@ protected AssistantConversationMessage() { AssistantConversationMessage( UUID id, UUID conversationId, + UUID turnId, UUID organizationId, UUID actorUserId, AssistantConversationRole role, @@ -49,6 +54,7 @@ protected AssistantConversationMessage() { Instant occurredAt) { super(Objects.requireNonNull(id, "id")); this.conversationId = Objects.requireNonNull(conversationId, "conversationId"); + this.turnId = Objects.requireNonNull(turnId, "turnId"); this.organizationId = Objects.requireNonNull(organizationId, "organizationId"); this.actorUserId = Objects.requireNonNull(actorUserId, "actorUserId"); this.role = Objects.requireNonNull(role, "role"); @@ -56,6 +62,10 @@ protected AssistantConversationMessage() { this.occurredAt = Objects.requireNonNull(occurredAt, "occurredAt"); } + UUID turnId() { + return turnId; + } + AssistantConversationMessageView view() { return new AssistantConversationMessageView( getId(), role, content, sequenceId, occurredAt, null); diff --git a/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationService.java b/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationService.java index 4651b525..66e1b4c7 100644 --- a/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationService.java +++ b/core/src/main/java/com/orgmemory/core/assistant/AssistantConversationService.java @@ -37,12 +37,13 @@ public class AssistantConversationService { } @Transactional - public UUID beginTurn(CurrentActor actor, UUID requestedId, String userMessage) { + public AssistantTurnRef beginTurn( + CurrentActor actor, UUID requestedId, String userMessage) { return beginTurn(actor, requestedId, userMessage, null); } @Transactional - public UUID beginTurn( + public AssistantTurnRef beginTurn( CurrentActor actor, UUID requestedId, String userMessage, @@ -62,15 +63,17 @@ public UUID beginTurn( conversation.touch(now); } conversation.selectModel(modelSelection); + UUID turnId = UUID.randomUUID(); messages.save(new AssistantConversationMessage( UUID.randomUUID(), conversation.getId(), + turnId, actor.organizationId(), actor.userId(), AssistantConversationRole.USER, validUserMessage, now)); - return conversation.getId(); + return new AssistantTurnRef(conversation.getId(), turnId); } @Transactional(readOnly = true) @@ -91,16 +94,16 @@ public void selectModel( @Transactional public void completeTurn( CurrentActor actor, - UUID conversationId, + AssistantTurnRef turn, UUID assistantMessageId, String assistantMessage) { - completeTurn(actor, conversationId, assistantMessageId, assistantMessage, List.of()); + completeTurn(actor, turn, assistantMessageId, assistantMessage, List.of()); } @Transactional public void completeTurn( CurrentActor actor, - UUID conversationId, + AssistantTurnRef turn, UUID assistantMessageId, String assistantMessage, List answerCitations) { @@ -108,11 +111,13 @@ public void completeTurn( return; } List persistedCitations = validateCitations(answerCitations); + UUID conversationId = turn.conversationId(); AssistantConversation conversation = requireOwned(actor, conversationId); Instant now = clock.instant(); messages.save(new AssistantConversationMessage( assistantMessageId, conversationId, + turn.turnId(), actor.organizationId(), actor.userId(), AssistantConversationRole.ASSISTANT, diff --git a/core/src/main/java/com/orgmemory/core/assistant/AssistantTurnRef.java b/core/src/main/java/com/orgmemory/core/assistant/AssistantTurnRef.java new file mode 100644 index 00000000..8725f4f8 --- /dev/null +++ b/core/src/main/java/com/orgmemory/core/assistant/AssistantTurnRef.java @@ -0,0 +1,21 @@ +package com.orgmemory.core.assistant; + +import java.util.Objects; +import java.util.UUID; + +/** + * Identifies one Assistant turn and the conversation it belongs to. + * + *

A turn spans two transactions — {@code beginTurn} persists the question and + * {@code completeTurn} persists the answer once the stream finishes. Passing this + * reference between them records the pairing the writers already know, so a + * reader never has to infer it from sequence order. Concurrent turns in one + * conversation can interleave their rows arbitrarily. + */ +public record AssistantTurnRef(UUID conversationId, UUID turnId) { + + public AssistantTurnRef { + Objects.requireNonNull(conversationId, "conversationId"); + Objects.requireNonNull(turnId, "turnId"); + } +} diff --git a/core/src/main/resources/db/migration/V26__assistant_message_turn_identity.sql b/core/src/main/resources/db/migration/V26__assistant_message_turn_identity.sql new file mode 100644 index 00000000..85416853 --- /dev/null +++ b/core/src/main/resources/db/migration/V26__assistant_message_turn_identity.sql @@ -0,0 +1,20 @@ +-- Give an Assistant turn an explicit identity. +-- +-- A turn writes its USER row in beginTurn and its ASSISTANT row in completeTurn, +-- in separate transactions. Two concurrent turns in one conversation can +-- therefore persist as U1, U2, A2, A1, which no ordering heuristic over +-- sequence_id pairs correctly. turn_id records the pairing the writers already +-- know instead of asking a later reader to infer it. +-- +-- Rows written before this migration stay NULL: they remain visible in the +-- product transcript and are ineligible as model context, because their pairing +-- cannot be recovered after the fact. + +ALTER TABLE public.assistant_conversation_messages + ADD COLUMN turn_id uuid; + +-- One USER and one ASSISTANT per turn. Partial so the legacy NULL rows, which +-- carry no pairing to protect, are exempt rather than colliding with each other. +CREATE UNIQUE INDEX idx_assistant_conversation_message_turn_role + ON public.assistant_conversation_messages (turn_id, role) + WHERE turn_id IS NOT NULL; diff --git a/core/src/test/java/com/orgmemory/core/assistant/AssistantConversationServiceTests.java b/core/src/test/java/com/orgmemory/core/assistant/AssistantConversationServiceTests.java index 9e53e59e..efc6a6cf 100644 --- a/core/src/test/java/com/orgmemory/core/assistant/AssistantConversationServiceTests.java +++ b/core/src/test/java/com/orgmemory/core/assistant/AssistantConversationServiceTests.java @@ -1,11 +1,13 @@ package com.orgmemory.core.assistant; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -58,10 +60,11 @@ void setUp() { @Test void createsAnActorOwnedConversationAndStoresTheRawUserTurn() { - UUID conversationId = service.beginTurn( + AssistantTurnRef turn = service.beginTurn( actor, null, " How do I submit an expense claim? "); + UUID conversationId = turn.conversationId(); ArgumentCaptor conversation = ArgumentCaptor.forClass(AssistantConversation.class); @@ -80,6 +83,47 @@ void createsAnActorOwnedConversationAndStoresTheRawUserTurn() { assertEquals( "How do I submit an expense claim?", message.getValue().view().content()); + assertEquals(turn.turnId(), message.getValue().turnId()); + } + + @Test + void writesBothHalvesOfOneTurnUnderTheIdentityBeginTurnAllocated() { + AssistantTurnRef turn = service.beginTurn(actor, null, "How long is probation?"); + when(conversations.findByIdAndOrganizationIdAndActorUserId( + turn.conversationId(), + actor.organizationId(), + actor.userId())) + .thenReturn(Optional.of(ownedConversation(turn.conversationId()))); + + service.completeTurn(actor, turn, UUID.randomUUID(), "Sixty days."); + + ArgumentCaptor saved = + ArgumentCaptor.forClass(AssistantConversationMessage.class); + verify(messages, times(2)).save(saved.capture()); + assertEquals( + List.of(AssistantConversationRole.USER, AssistantConversationRole.ASSISTANT), + saved.getAllValues().stream().map(m -> m.view().role()).toList()); + // The pairing is recorded by the writers, not inferred later from + // sequence order, which concurrent turns interleave. + assertEquals( + List.of(turn.turnId(), turn.turnId()), + saved.getAllValues().stream().map(AssistantConversationMessage::turnId).toList()); + } + + @Test + void givesConcurrentTurnsInOneConversationDistinctIdentities() { + AssistantConversation conversation = ownedConversation(UUID.randomUUID()); + when(conversations.findForUpdateByIdAndOrganizationIdAndActorUserId( + conversation.getId(), + actor.organizationId(), + actor.userId())) + .thenReturn(Optional.of(conversation)); + + AssistantTurnRef first = service.beginTurn(actor, conversation.getId(), "First"); + AssistantTurnRef second = service.beginTurn(actor, conversation.getId(), "Second"); + + assertEquals(first.conversationId(), second.conversationId()); + assertNotEquals(first.turnId(), second.turnId()); } @Test @@ -144,7 +188,7 @@ void persistsTheServerAllocatedAssistantMessageIdentity() { service.completeTurn( actor, - conversationId, + new AssistantTurnRef(conversationId, UUID.randomUUID()), messageId, "The probation period is 60 days. [1]"); @@ -172,7 +216,7 @@ void persistsServerDeclaredCitationReferencesWithTheCompletedAnswer() { service.completeTurn( actor, - conversationId, + new AssistantTurnRef(conversationId, UUID.randomUUID()), messageId, "The probation period is 60 days. [1]", List.of(new AssistantCitation(1, evidence))); @@ -193,6 +237,7 @@ void readsCitationReferencesOnlyThroughNonLockingActorOwnedAssistantLookup() { AssistantConversationMessage answer = new AssistantConversationMessage( messageId, conversationId, + UUID.randomUUID(), actor.organizationId(), actor.userId(), AssistantConversationRole.ASSISTANT, @@ -249,6 +294,7 @@ void createsReplacesAndRemovesFeedbackForAnOwnedAssistantAnswer() { AssistantConversationMessage answer = new AssistantConversationMessage( messageId, conversationId, + UUID.randomUUID(), actor.organizationId(), actor.userId(), AssistantConversationRole.ASSISTANT, @@ -321,9 +367,11 @@ void rejectsAnInvalidConversationTitleAsBusinessValidation() { void returnsTheFullOwnedTranscriptInPersistedOrder() { UUID conversationId = UUID.randomUUID(); AssistantConversation conversation = ownedConversation(conversationId); + UUID turnId = UUID.randomUUID(); AssistantConversationMessage first = new AssistantConversationMessage( UUID.randomUUID(), conversationId, + turnId, actor.organizationId(), actor.userId(), AssistantConversationRole.USER, @@ -332,6 +380,7 @@ void returnsTheFullOwnedTranscriptInPersistedOrder() { AssistantConversationMessage second = new AssistantConversationMessage( UUID.randomUUID(), conversationId, + turnId, actor.organizationId(), actor.userId(), AssistantConversationRole.ASSISTANT, diff --git a/docs/increments/active/2026-08-06-assistant-conversation-memory-ssot/plan.md b/docs/increments/active/2026-08-06-assistant-conversation-memory-ssot/plan.md index 356392c3..6625b2f5 100644 --- a/docs/increments/active/2026-08-06-assistant-conversation-memory-ssot/plan.md +++ b/docs/increments/active/2026-08-06-assistant-conversation-memory-ssot/plan.md @@ -6,10 +6,10 @@ with a no-tools judge; record brief and verdict. - [x] Port status-mapped failure sentences so a failed turn names its cause. Independent of the store change and shippable on its own. -- [ ] Add `turn_id` to `assistant_conversation_messages` with partial uniqueness +- [x] Add `turn_id` to `assistant_conversation_messages` with partial uniqueness for one USER and one ASSISTANT per turn; leave legacy rows nullable, transcript-visible and context-ineligible. -- [ ] Carry the turn id through `beginTurn` and `completeTurn` so the pair is +- [x] Carry the turn id through `beginTurn` and `completeTurn` so the pair is written against one identity rather than inferred from sequence order. - [ ] Add a project-owned read-only transcript context advisor that selects the last completed turns, excludes the in-flight USER by construction, and snaps diff --git a/docs/specs/domains/assistant-and-mcp.md b/docs/specs/domains/assistant-and-mcp.md index 325665b6..117127ba 100644 --- a/docs/specs/domains/assistant-and-mcp.md +++ b/docs/specs/domains/assistant-and-mcp.md @@ -9,7 +9,7 @@ Source: `core/src/main/java/com/orgmemory/core/assistant`, `apps/web/src/features/assistant`, and `apps/web/src/components/ai-elements/model-selector.tsx`. -Reconciled: `2026-08-06-assistant-skill-activity-receipt (64221f86)`. +Reconciled: `2026-08-06-assistant-turn-identity (e13685eb)`. ## Current Behavior @@ -195,6 +195,15 @@ recheck current visibility in fixed batches of at most 20. Determinate denied, missing, stale, or revised sources simply leave their historical marker inert; an indeterminate authorization result fails citation hydration without hiding the saved answer. + +Each turn carries an explicit identity allocated when its question is persisted +and carried through to its answer, so both rows record the pairing their writers +already knew. A turn holds at most one question and one answer, enforced by a +partial unique index over turn and role. Two turns in one conversation may both +open before either answers and may persist their rows in any interleaving, so +sequence order does not pair them. Rows written before turn identity existed +carry none; they remain in the transcript and cannot be paired. + Each conversation also stores its optional model activation together with the organization route override identity and version observed at selection. Picker changes and turn creation lock the owned conversation row. Disabled catalog diff --git a/docs/tests/domains/assistant-and-mcp.md b/docs/tests/domains/assistant-and-mcp.md index 77454299..f23dfb4d 100644 --- a/docs/tests/domains/assistant-and-mcp.md +++ b/docs/tests/domains/assistant-and-mcp.md @@ -8,7 +8,7 @@ Source: `core/src/test/java/com/orgmemory/core/assistant`, `apps/web/src/features/assistant`, plus `apps/web/test/e2e/assistant-pipeline.spec.ts`. -Reconciled: `2026-08-06-assistant-skill-activity-receipt (64221f86)`. +Reconciled: `2026-08-06-assistant-turn-identity (e13685eb)`. | Behavior | Evidence | Status | | --- | --- | --- | @@ -53,6 +53,7 @@ Reconciled: `2026-08-06-assistant-skill-activity-receipt (64221f86)`. | Asset tool traces contain exact release refs without raw Prompt secrets/output | `AssistantAssetToolServiceTests#promptTraceStoresShapeAndDigestButNoRawSecretOrOutput` | covered | | Asset Assistant has no approval/publication/withdrawal/permission/arbitrary-execution action | `AssistantAssetToolServiceTests#assistantActionRegistryHasNoGovernanceOrArbitraryExecutionPath` | covered | | Full transcript is actor-owned, replayed in order, and rejects another actor before writing | `AssistantConversationServiceTests` | covered | +| A turn's question and answer are paired by an explicit identity, hold at most one row per role, survive out-of-order persistence, and leave pre-existing rows unpaired | `AssistantTurnIdentityIntegrationTests`, `AssistantConversationServiceTests#writesBothHalvesOfOneTurnUnderTheIdentityBeginTurnAllocated`, `AssistantConversationServiceTests#givesConcurrentTurnsInOneConversationDistinctIdentities` | covered | | Completed answers atomically persist ordered citation mappings with composite ownership, uniqueness, and cascade deletion | `AssistantConversationServiceTests#persistsServerDeclaredCitationReferencesWithTheCompletedAnswer`, `AssistantMessageCitationMigrationTests` | covered | | Citation hydration is transcript-independent, reloadable, actor-owned, bounded to 100, deduplicated, and current-authorization filtered | `AssistantConversationServiceTests` citation-reference scenarios, `CanonicalEvidenceAuthorizationServiceTests`, `CitationEvidenceServiceTests`, `assistant-pipeline.spec.ts#rehydrates currently authorized citations after transcript reload` | covered | | Excerpts reauthorize current evidence, cap Unicode content, audit allow/deny, and keep missing/revoked/stale outcomes opaque | `CitationEvidenceServiceTests`, `CitationContentControllerTests`, `CitationContentWebMvcTests`, `assistant-pipeline.spec.ts` revocation scenario | covered | @@ -85,3 +86,4 @@ Reconciled: `2026-08-06-assistant-skill-activity-receipt (64221f86)`. | Every row above about time to first token holds against the handler in isolation and held while the handler was never registered in a running application; only `ConfigurationConditionTests` and production traffic distinguish the two | `AssistantTurnObservationTests` construct the handler directly | partial | | Every meter a dashboard charts as a quantile publishes a bounded percentile histogram | `MetricsDistributionTests` | covered | | General chat-turn idempotency | none | not implemented | +| Two turns of one conversation completing at the same instant | none | gap — `completeTurn` touches the conversation without a lock, so simultaneous completion raises an optimistic-locking failure and loses one answer; pre-existing and unrelated to turn identity |