Files
twenty/packages/twenty-front/src/modules/ai/components/AgentChatMessagesFetchEffect.tsx
T
Félix Malfait 77529191f8 fix(ai): apply stream chunks in exact server seq order via a client-side sequencer (#22484)
## Rationale

Stream chunks reach the client on two unsynchronized paths: live SSE
events and the catchup replay (fired on reload, refetch, SSE reconnect,
and keep-alive recovery). The server already stamps every chunk with an
authoritative `seq` (Redis `RPUSH` length), but the client applies
chunks in **arrival order**. Reload mid-stream and the two paths
interleave: duplicated text deltas, or lower-seq catchup chunks applied
after higher-seq live ones — the streaming answer visibly garbles until
the persist-refetch repaints it.

Main's existing guard (`seq < firstLiveSeq` bound on catchup) only
prevents duplication in one direction (live-before-catchup); it does
nothing for catchup-during-live overlap, and it *creates* a
dropped-chunk window when chunks land between the catchup snapshot and
the first live event.

## Why this is the root cause, not a symptom patch

The defect is a joining problem between two ordered sources, and the
join point is the client — the server can't fix it without a protocol
change (per-subscriber cursor resume), because Redis pub/sub fan-out has
no per-subscriber replay. Given the transport, the correct fix is to
make the reducer's input **seq-exact**: apply strictly in server order,
dedup anything already applied, buffer early arrivals until the gap
fills. Escalation is bounded and degrades gracefully: a stalled gap
triggers one refetch (the full-list catchup replay doubles as gap-fill,
no new endpoint), a second stall flushes the buffer in order — so even
an expired chunk list degrades to slightly-lossy instead of wedging. The
catchup path now replays the full list (the sequencer dedups overlap),
which also closes the dropped-chunk window.

Server-side cursor resume remains the nicer long-term protocol (would
simplify this client), but it's a subscription protocol change; this
fixes the user-facing defect with zero server change and is
forward-compatible with it.

## User impact

Reloading (or losing the connection) mid-answer currently scrambles or
duplicates the streaming text until the turn completes. With this, the
answer renders identically no matter when you reload or how the two
delivery paths race.

## Test plan

- [x] Sequencer unit suite (fake timers): in-order apply, out-of-order
buffering, catchup/live overlap dedup, gap-fill via replay,
stall→refetch escalation, second-stall in-order flush, high-water-mark
continuation, reset
- [ ] CI green
- [ ] Manual: reload mid-stream repeatedly; text never reorders

https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38

---
_Generated by [Claude
Code](https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38)_

<!-- This is an auto-generated description by cubic. -->
<a
href="https://cubic.dev/pr/twentyhq/twenty/pull/22484?utm_source=github"
target="_blank" rel="noopener noreferrer"
data-no-image-dialog="true"><picture><source
media="(prefers-color-scheme: dark)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"><source
media="(prefers-color-scheme: light)"
srcset="https://www.cubic.dev/buttons/review-in-cubic-light.svg"><img
alt="Review in cubic"
src="https://www.cubic.dev/buttons/review-in-cubic-dark.svg"></picture></a>
<!-- End of auto-generated description by cubic. -->
2026-07-02 21:17:11 +02:00

181 lines
6.3 KiB
TypeScript

import { useStore } from 'jotai';
import { useCallback, useMemo } from 'react';
import { type AgentChatSubscriptionEvent } from 'twenty-shared/ai';
import { isDefined } from 'twenty-shared/utils';
import { AGENT_CHAT_REFETCH_MESSAGES_EVENT_NAME } from '@/ai/constants/AgentChatRefetchMessagesEventName';
import { AGENT_CHAT_NEW_THREAD_DRAFT_KEY } from '@/ai/states/agentChatDraftsByThreadIdState';
import { agentChatFetchedMessagesComponentFamilyState } from '@/ai/states/agentChatFetchedMessagesComponentFamilyState';
import { agentChatFirstLiveSeqComponentFamilyState } from '@/ai/states/agentChatFirstLiveSeqComponentFamilyState';
import { agentChatHandleEventCallbackComponentFamilyState } from '@/ai/states/agentChatHandleEventCallbackComponentFamilyState';
import { agentChatIsAwaitingPersistedRefetchComponentFamilyState } from '@/ai/states/agentChatIsAwaitingPersistedRefetchComponentFamilyState';
import { agentChatMessagesLoadingState } from '@/ai/states/agentChatMessagesLoadingState';
import { agentChatQueuedMessagesComponentFamilyState } from '@/ai/states/agentChatQueuedMessagesComponentFamilyState';
import { currentAiChatThreadState } from '@/ai/states/currentAiChatThreadState';
import { skipMessagesSkeletonUntilLoadedState } from '@/ai/states/skipMessagesSkeletonUntilLoadedState';
import { mapDBMessagesToUIMessages } from '@/ai/utils/mapDBMessagesToUIMessages';
import { SSE_CLIENT_RECONNECTED_EVENT_NAME } from '@/sse-db-event/constants/SseClientReconnectedEventName';
import { useQueryWithCallbacks } from '@/apollo/hooks/useQueryWithCallbacks';
import { useListenToBrowserEvent } from '@/browser-event/hooks/useListenToBrowserEvent';
import { useAtomComponentFamilyStateCallbackState } from '@/ui/utilities/state/jotai/hooks/useAtomComponentFamilyStateCallbackState';
import { useAtomStateValue } from '@/ui/utilities/state/jotai/hooks/useAtomStateValue';
import { useSetAtomComponentFamilyState } from '@/ui/utilities/state/jotai/hooks/useSetAtomComponentFamilyState';
import { useSetAtomState } from '@/ui/utilities/state/jotai/hooks/useSetAtomState';
import {
GetChatMessagesDocument,
type GetChatMessagesQuery,
} from '~/generated-metadata/graphql';
export const AgentChatMessagesFetchEffect = () => {
const store = useStore();
const currentAiChatThread = useAtomStateValue(currentAiChatThreadState);
const isNewThread = useMemo(
() =>
currentAiChatThread === null ||
currentAiChatThread === AGENT_CHAT_NEW_THREAD_DRAFT_KEY,
[currentAiChatThread],
);
const setAgentChatMessagesLoading = useSetAtomState(
agentChatMessagesLoadingState,
);
const setSkipMessagesSkeletonUntilLoaded = useSetAtomState(
skipMessagesSkeletonUntilLoadedState,
);
const setAgentChatFetchedMessages = useSetAtomComponentFamilyState(
agentChatFetchedMessagesComponentFamilyState,
{ threadId: currentAiChatThread },
);
const setAgentChatQueuedMessages = useSetAtomComponentFamilyState(
agentChatQueuedMessagesComponentFamilyState,
{ threadId: currentAiChatThread },
);
const setAgentChatIsAwaitingPersistedRefetch = useSetAtomComponentFamilyState(
agentChatIsAwaitingPersistedRefetchComponentFamilyState,
{ threadId: currentAiChatThread },
);
const handleEventCallbackFamilyCallback =
useAtomComponentFamilyStateCallbackState(
agentChatHandleEventCallbackComponentFamilyState,
);
const firstLiveSeqFamilyCallback = useAtomComponentFamilyStateCallbackState(
agentChatFirstLiveSeqComponentFamilyState,
);
const handleFirstLoad = useCallback(
(_data: GetChatMessagesQuery) => {
setSkipMessagesSkeletonUntilLoaded(false);
},
[setSkipMessagesSkeletonUntilLoaded],
);
const handleDataLoaded = useCallback(
(data: GetChatMessagesQuery) => {
const uiMessages = mapDBMessagesToUIMessages(data.chatMessages ?? []);
setAgentChatFetchedMessages(
uiMessages.filter((message) => message.status !== 'queued'),
);
setAgentChatQueuedMessages(
uiMessages.filter((message) => message.status === 'queued'),
);
setAgentChatIsAwaitingPersistedRefetch(false);
const catchup = data.chatStreamCatchupChunks;
if (!isDefined(catchup)) {
return;
}
const threadId = store.get(currentAiChatThreadState.atom);
if (!isDefined(threadId)) {
return;
}
const familyKey = { threadId };
const handleEvent = store.get(
handleEventCallbackFamilyCallback(familyKey),
);
if (!isDefined(handleEvent)) {
return;
}
const firstLiveSeq = store.get(firstLiveSeqFamilyCallback(familyKey));
for (let index = 0; index < catchup.chunks.length; index++) {
handleEvent({
type: 'stream-chunk',
chunk: catchup.chunks[index],
seq: index + 1,
} as AgentChatSubscriptionEvent);
}
if (isDefined(catchup.error) && firstLiveSeq === null) {
handleEvent({
type: 'stream-error',
code: catchup.error.code,
message: catchup.error.message,
} as AgentChatSubscriptionEvent);
}
},
[
setAgentChatFetchedMessages,
setAgentChatQueuedMessages,
setAgentChatIsAwaitingPersistedRefetch,
store,
handleEventCallbackFamilyCallback,
firstLiveSeqFamilyCallback,
],
);
const handleLoadingChange = useCallback(
(loading: boolean) => {
setAgentChatMessagesLoading(loading);
if (!loading) {
setAgentChatIsAwaitingPersistedRefetch(false);
}
},
[setAgentChatMessagesLoading, setAgentChatIsAwaitingPersistedRefetch],
);
const { refetch: refetchAgentChatMessages } = useQueryWithCallbacks(
GetChatMessagesDocument,
{
variables: { threadId: currentAiChatThread ?? '' },
skip: !isDefined(currentAiChatThread) || isNewThread,
onFirstLoad: handleFirstLoad,
onDataLoaded: handleDataLoaded,
onLoadingChange: handleLoadingChange,
},
);
const handleRefetchMessages = useCallback(() => {
if (isNewThread) {
return;
}
refetchAgentChatMessages();
}, [refetchAgentChatMessages, isNewThread]);
useListenToBrowserEvent({
eventName: AGENT_CHAT_REFETCH_MESSAGES_EVENT_NAME,
onBrowserEvent: handleRefetchMessages,
});
useListenToBrowserEvent({
eventName: SSE_CLIENT_RECONNECTED_EVENT_NAME,
onBrowserEvent: handleRefetchMessages,
});
return null;
};