fix(api): order service timeouts so inner hops expire first (#2131)

This commit is contained in:
Hampus
2026-08-30 22:15:36 +02:00
committed by GitHub
parent 9d1733bdb3
commit e14d193b43
5 changed files with 21 additions and 5 deletions
@@ -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<Record<string, unknown>> = [];
readonly timeouts: Array<number | undefined> = [];
async connect(): Promise<void> {}
@@ -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<string, unknown>);
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);
});
});
@@ -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;
@@ -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);
@@ -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;
@@ -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<string>, fallback: number): number {
const envNames = Array.isArray(names) ? names : [names];