From ca8ab32253879227c995998c54fa17c013c58ae9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?F=C3=A9lix=20Malfait?= Date: Thu, 2 Jul 2026 21:05:19 +0200 Subject: [PATCH] fix(ai): delete the Redis chunk list on credits-exhausted terminal events (#22477) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Rationale The AI chat stream keeps every published chunk in a Redis list (`agent-chat-stream-chunks:`, 1h TTL) so late subscribers can catch up. On `message-persisted` the list is deleted. On `credits-exhausted` — the *other* successful terminal event — it wasn't. Any reload/refetch within the TTL replayed the orphaned chunks, flipping the thread into a "streaming" state that no terminal event ever closes: an endless spinner until the user sends another message. **Production evidence (Sentry):** `Billing Credits Exhausted` fired for **290 users / 937 events in 90 days**, ongoing ([TWENTY-SERVER-G42](https://twenty-v7.sentry.io/issues/TWENTY-SERVER-G42)) — every one of those users who reloads the chat within an hour hits this. ## Why this is the root cause, not a symptom patch The chunk list's lifecycle contract is "cleared when the turn settles". `credits-exhausted` resolves the job successfully **without** persisting a `lastStreamError`, so unlike `stream-error` there is no persisted terminator for catchup to replay after the chunks — the replay is unconditionally un-closeable. Deleting on both settle events restores the contract exactly where it's already enforced for `message-persisted`. `stream-error` deliberately keeps the list: the persisted error acts as the replay terminator, letting a reloading client still see the failed turn's partial output. ## User impact Users who hit their billing cap mid-answer (~100/month) no longer come back to a permanently spinning thread after a reload — they see the settled conversation and the billing state. ## Test plan - [x] Unit spec: chunk accumulation with 1-based seq, deletion on both terminal events (`it.each`), retention on `stream-error` - [ ] CI green https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38 --- _Generated by [Claude Code](https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38)_ Review in cubic --- ...agent-chat-event-publisher.service.spec.ts | 76 +++++++++++++++++++ .../agent-chat-event-publisher.service.ts | 5 +- 2 files changed, 80 insertions(+), 1 deletion(-) create mode 100644 packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-event-publisher.service.spec.ts diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-event-publisher.service.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-event-publisher.service.spec.ts new file mode 100644 index 0000000000..19e22aad8e --- /dev/null +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-event-publisher.service.spec.ts @@ -0,0 +1,76 @@ +import { AgentChatEventPublisherService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service'; + +describe('AgentChatEventPublisherService', () => { + const buildService = () => { + const redis = { + rpush: jest.fn().mockResolvedValue(1), + expire: jest.fn().mockResolvedValue(1), + del: jest.fn().mockResolvedValue(1), + }; + const subscriptionService = { + publishToAgentChat: jest.fn().mockResolvedValue(undefined), + }; + const service = new AgentChatEventPublisherService( + subscriptionService as never, + { getClient: () => redis } as never, + ); + + return { service, redis, subscriptionService }; + }; + + it('accumulates stream chunks with a 1-based sequence number', async () => { + const { service, redis, subscriptionService } = buildService(); + + redis.rpush.mockResolvedValue(7); + + await service.publish({ + threadId: 'thread-id', + workspaceId: 'workspace-id', + event: { type: 'stream-chunk', chunk: { type: 'text-delta' } } as never, + }); + + expect(redis.rpush).toHaveBeenCalledWith( + 'agent-chat-stream-chunks:thread-id', + JSON.stringify({ type: 'text-delta' }), + ); + expect(subscriptionService.publishToAgentChat).toHaveBeenCalledWith( + expect.objectContaining({ + payload: { + onAgentChatEvent: { + threadId: 'thread-id', + event: expect.objectContaining({ seq: 7 }), + }, + }, + }), + ); + }); + + it.each(['message-persisted', 'credits-exhausted'] as const)( + 'deletes the accumulated chunk list on terminal %s', + async (type) => { + const { service, redis } = buildService(); + + await service.publish({ + threadId: 'thread-id', + workspaceId: 'workspace-id', + event: { type } as never, + }); + + expect(redis.del).toHaveBeenCalledWith( + 'agent-chat-stream-chunks:thread-id', + ); + }, + ); + + it('keeps the chunk list on stream-error so catchup can replay the partial turn', async () => { + const { service, redis } = buildService(); + + await service.publish({ + threadId: 'thread-id', + workspaceId: 'workspace-id', + event: { type: 'stream-error', code: 'X', message: 'boom' } as never, + }); + + expect(redis.del).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service.ts index a666187975..aec2cc8d17 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service.ts @@ -38,7 +38,10 @@ export class AgentChatEventPublisherService { await redis.expire(key, STREAM_CHUNKS_TTL_SECONDS); publishedEvent = { ...event, seq }; - } else if (event.type === 'message-persisted') { + } else if ( + event.type === 'message-persisted' || + event.type === 'credits-exhausted' + ) { const redis = this.redisClientService.getClient(); await redis.del(this.getStreamChunksKey(threadId)); }