From 1feb41eb6483238f0ab104df896333ee3c30d844 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?F=C3=A9lix=20Malfait?= Date: Thu, 2 Jul 2026 21:11:27 +0200 Subject: [PATCH] perf(sse): drop the per-workspace lock around atomic activeStreams set operations (#22476) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Rationale Every event-stream create/destroy in a workspace serializes on a single cache lock (`workspace::activeStreams`) just to run `setAdd`/`setRemove`. On Redis those are native `SADD`/`SREM` — already atomic (`cache-storage.service.ts:79-109`) — so the lock adds zero correctness. What it does add: 100ms lock-retry polling under concurrency, and a hard 5s ceiling (50 retries × 100ms, `cache-lock.service.ts:36`) after which stream creation **fails** with `Failed to acquire lock`. **Production evidence (Sentry):** [TWENTY-SERVER-H3W](https://twenty-v7.sentry.io/issues/TWENTY-SERVER-H3W) (125 events, 16 users) server-side and [TWENTY-FRONT-6MM](https://twenty-v7.sentry.io/issues/TWENTY-FRONT-6MM) (65 users) client-side — spiky, consistent with deploy-driven reconnect storms in large workspaces where every tab reconnects at once and queues on one lock. ## Why this is the root cause, not a symptom patch The failure isn't "the lock timeout is too short" (raising it would just trade errors for latency) — it's that the critical section doesn't exist. A set-membership add/remove of a single element has no read-modify-write window on Redis. The lock that **is** legitimate stays untouched: `@WithLock` on `addQuery`/`removeQuery`, which genuinely read-modify-write the JSON queries map. The non-Redis fallback of `setAdd` is read-modify-write, but that path only serves single-node dev/test setups where the worst case is a transiently miscounted metrics gauge (`twenty_event_streams_live_total`), not a correctness issue — the authoritative per-stream state lives in its own key. ## User impact During reconnect storms (deploys, network blips) in busy workspaces, tabs no longer randomly fail to establish their event stream — which previously meant no live updates for that tab and, through the strict client error path, a crash of the listener sync loop (fixed separately in #22475). Also removes up-to-5s of serialized queueing latency per workspace on every connect wave. ## Test plan - [x] Verified `setAdd`/`setRemove` are native `SADD`/`SREM` + `EXPIRE` on the Redis path - [x] No callers depend on the lock's ordering (checked `engine/subscriptions` call sites) - [ ] CI green https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38 --- _Generated by [Claude Code](https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38)_ Review in cubic --- .../subscriptions/event-stream.service.ts | 39 +++++++------------ 1 file changed, 13 insertions(+), 26 deletions(-) diff --git a/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts b/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts index 5dbe44d9f0..d043d255b5 100644 --- a/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts +++ b/packages/twenty-server/src/engine/subscriptions/event-stream.service.ts @@ -3,7 +3,6 @@ import { Injectable, Logger, OnModuleInit } from '@nestjs/common'; import { isDefined } from 'twenty-shared/utils'; import { type SerializableAuthContext } from 'src/engine/core-modules/auth/types/serializable-auth-context.type'; -import { CacheLockService } from 'src/engine/core-modules/cache-lock/cache-lock.service'; import { WithLock } from 'src/engine/core-modules/cache-lock/with-lock.decorator'; import { InjectCacheStorage } from 'src/engine/core-modules/cache-storage/decorators/cache-storage.decorator'; import { CacheStorageService } from 'src/engine/core-modules/cache-storage/services/cache-storage.service'; @@ -30,7 +29,6 @@ export class EventStreamService implements OnModuleInit { constructor( @InjectCacheStorage(CacheStorageNamespace.EngineSubscriptions) private readonly cacheStorageService: CacheStorageService, - private readonly cacheLockService: CacheLockService, private readonly metricsService: MetricsService, ) {} @@ -90,15 +88,11 @@ export class EventStreamService implements OnModuleInit { await this.cacheStorageService.set(key, streamData, EVENT_STREAM_TTL_MS); - const activeStreamsKey = this.getActiveStreamsKey(workspaceId); - - await this.cacheLockService.withLock(async () => { - await this.cacheStorageService.setAdd( - activeStreamsKey, - [eventStreamChannelId], - EVENT_STREAM_TTL_MS, - ); - }, activeStreamsKey); + await this.cacheStorageService.setAdd( + this.getActiveStreamsKey(workspaceId), + [eventStreamChannelId], + EVENT_STREAM_TTL_MS, + ); } async destroyEventStream({ @@ -112,13 +106,10 @@ export class EventStreamService implements OnModuleInit { await this.cacheStorageService.del(key); - const activeStreamsKey = this.getActiveStreamsKey(workspaceId); - - await this.cacheLockService.withLock(async () => { - await this.cacheStorageService.setRemove(activeStreamsKey, [ - eventStreamChannelId, - ]); - }, activeStreamsKey); + await this.cacheStorageService.setRemove( + this.getActiveStreamsKey(workspaceId), + [eventStreamChannelId], + ); } async getActiveStreamIds(workspaceId: string): Promise { @@ -135,14 +126,10 @@ export class EventStreamService implements OnModuleInit { return; } - const activeStreamsKey = this.getActiveStreamsKey(workspaceId); - - await this.cacheLockService.withLock(async () => { - await this.cacheStorageService.setRemove( - activeStreamsKey, - streamIdsToRemove, - ); - }, activeStreamsKey); + await this.cacheStorageService.setRemove( + this.getActiveStreamsKey(workspaceId), + streamIdsToRemove, + ); } async getStreamsData(