diff --git a/fluxer_api/src/api/channel/services/message/MessageResponseDataService.test.ts b/fluxer_api/src/api/channel/services/message/MessageResponseDataService.test.ts index 1c1a2b755..18dde62f1 100644 --- a/fluxer_api/src/api/channel/services/message/MessageResponseDataService.test.ts +++ b/fluxer_api/src/api/channel/services/message/MessageResponseDataService.test.ts @@ -10,9 +10,11 @@ import {MessageResponseDataService} from './MessageResponseDataService'; const encoder = new TextEncoder(); const decoder = new TextDecoder(); +const ROUTER_SHARD_REQUEST_TIMEOUT_MS = 5000; class FakeConnectionManager implements INatsConnectionManager { readonly payloads: Array> = []; + readonly timeouts: Array = []; async connect(): Promise {} @@ -24,8 +26,9 @@ class FakeConnectionManager implements INatsConnectionManager { getConnection(): NatsConnection { return { - request: async (_subject: string, data: Uint8Array) => { + request: async (_subject: string, data: Uint8Array, options?: {timeout?: number}) => { this.payloads.push(JSON.parse(decoder.decode(data)) as Record); + this.timeouts.push(options?.timeout); return { data: encoder.encode( JSON.stringify({ @@ -117,4 +120,17 @@ describe('MessageResponseDataService', () => { viewer_user_id: '3', }); }); + + it('waits longer than the router shard timeout so the inner hop expires first', async () => { + const connectionManager = new FakeConnectionManager(); + const service = new MessageResponseDataService(connectionManager); + + await service.buildMessageForChannel({ + channel: {guildId: null}, + message: makeMessage(), + }); + + expect(connectionManager.timeouts[0]).toBe(6000); + expect(connectionManager.timeouts[0]).toBeGreaterThan(ROUTER_SHARD_REQUEST_TIMEOUT_MS); + }); }); diff --git a/fluxer_api/src/api/channel/services/message/MessageResponseDataService.ts b/fluxer_api/src/api/channel/services/message/MessageResponseDataService.ts index bc41d00bf..b0cb265fc 100644 --- a/fluxer_api/src/api/channel/services/message/MessageResponseDataService.ts +++ b/fluxer_api/src/api/channel/services/message/MessageResponseDataService.ts @@ -13,7 +13,7 @@ import type {Message} from '../../../models/Message'; import {isJsonRecord, parseJsonWithGuard} from '../../../utils/JsonBoundaryUtils'; const MESSAGE_RESPONSE_SERVICE_SUBJECT = 'svc.messages'; -const MESSAGE_RESPONSE_SERVICE_TIMEOUT_MS = 3000; +const MESSAGE_RESPONSE_SERVICE_TIMEOUT_MS = 6000; let messageResponseDataService: MessageResponseDataService | undefined; let injectedMessageResponseDataService: MessageResponseDataService | undefined; diff --git a/fluxer_api/src/api/infrastructure/SnowflakeService.test.ts b/fluxer_api/src/api/infrastructure/SnowflakeService.test.ts index 594a0bc8b..69ea02aa8 100644 --- a/fluxer_api/src/api/infrastructure/SnowflakeService.test.ts +++ b/fluxer_api/src/api/infrastructure/SnowflakeService.test.ts @@ -116,7 +116,7 @@ describe('SnowflakeService', () => { count: 1, routing_key: 'channel:1510189013330296832', }, - timeout: 5000, + timeout: 6000, }); expect(await service.generate()).toBe(201n); expect(manager.requests).toHaveLength(2); diff --git a/fluxer_api/src/api/infrastructure/SnowflakeService.ts b/fluxer_api/src/api/infrastructure/SnowflakeService.ts index dbfce93ae..03d8a6d39 100644 --- a/fluxer_api/src/api/infrastructure/SnowflakeService.ts +++ b/fluxer_api/src/api/infrastructure/SnowflakeService.ts @@ -9,7 +9,7 @@ import type {ISnowflakeService} from './ISnowflakeService'; const DEFAULT_REMOTE_SUBJECT = 'svc.snowflakes'; const DEFAULT_REMOTE_BATCH_SIZE = 128; const DEFAULT_REMOTE_LOW_WATERMARK = 32; -const DEFAULT_REMOTE_TIMEOUT_MS = 5000; +const DEFAULT_REMOTE_TIMEOUT_MS = 6000; const DEFAULT_REMOTE_MAX_BUFFER_AGE_MS = 5000; const MAX_REMOTE_BATCH_SIZE = 512; diff --git a/fluxer_api/src/api/middleware/ServiceRegistry.ts b/fluxer_api/src/api/middleware/ServiceRegistry.ts index b80db9161..14d7722c0 100644 --- a/fluxer_api/src/api/middleware/ServiceRegistry.ts +++ b/fluxer_api/src/api/middleware/ServiceRegistry.ts @@ -34,7 +34,7 @@ import type {WorkerTaskName} from '../worker/WorkerLaneConfig'; const DEFAULT_SNOWFLAKE_SERVICE_BATCH_SIZE = 128; const DEFAULT_SNOWFLAKE_SERVICE_LOW_WATERMARK = 32; const DEFAULT_SNOWFLAKE_SERVICE_MAX_BUFFER_AGE_MS = 5000; -const DEFAULT_SNOWFLAKE_SERVICE_REQUEST_TIMEOUT_MS = 5000; +const DEFAULT_SNOWFLAKE_SERVICE_REQUEST_TIMEOUT_MS = 6000; function readPositiveIntegerEnv(names: string | Array, fallback: number): number { const envNames = Array.isArray(names) ? names : [names];