Таймауты отправки

При работе с брокерами сообщений отправка данных не всегда завершается мгновенно. Сетевые задержки, перегрузка брокера, нестабильное соединение, проблемы маршрутизации и ограничения TCP способны привести к зависанию операции передачи. В таких ситуациях приложение может бесконечно ожидать завершения отправки сообщения.

Таймауты отправки позволяют ограничить время ожидания операции публикации. Если сообщение не было передано за установленный промежуток времени, операция считается неуспешной.

В экосистеме STOMP поверх WebSocket тема таймаутов особенно важна, поскольку:

  • WebSocket может оставаться формально открытым даже при фактической потере соединения;
  • браузер не предоставляет полного контроля над TCP-состоянием;
  • подтверждение доставки зависит от брокера;
  • отсутствие heartbeat-механизма может скрывать разрыв соединения;
  • асинхронная отправка затрудняет определение момента ошибки.

Особенности отправки сообщений в STOMP.js

Метод publish() в STOMP.js не возвращает Promise и обычно выполняется мгновенно:

client.publish({
    destination: '/topic/chat',
    body: 'Сообщение'
});

На первый взгляд может показаться, что сообщение успешно отправлено сразу после вызова метода. На практике это означает лишь передачу данных в WebSocket-слой браузера.

Фактическая доставка зависит от:

  • состояния WebSocket;
  • доступности брокера;
  • наличия сетевого соединения;
  • внутренней очереди браузера;
  • настроек heartbeat;
  • подтверждений ACK/NACK.

Из-за этого в STOMP.js таймауты реализуются не как встроенная функция метода publish(), а через комбинацию дополнительных механизмов контроля.


Проблема зависающих отправок

Типичная проблема возникает при частичном разрыве сети.

Например:

  1. WebSocket ещё считается открытым;
  2. браузер продолжает принимать вызовы send;
  3. пакеты физически не доходят до брокера;
  4. приложение считает сообщения отправленными;
  5. пользователь теряет данные.

Такое состояние называется «half-open connection».

Без таймаутов приложение не способно обнаружить проблему своевременно.


Контроль через heartbeat

Наиболее распространённый способ обнаружения проблем — heartbeat.

В STOMP.js heartbeat настраивается так:

import { Client } from '@stomp/stompjs';

const client = new Client({
    brokerURL: 'ws://localhost:15674/ws',

    heartbeatIncoming: 10000,
    heartbeatOutgoing: 10000
});

Значение параметров

heartbeatOutgoing

Интервал отправки heartbeat-пакетов серверу.

heartbeatIncoming

Максимальное время ожидания heartbeat от брокера.


Как heartbeat связан с таймаутами отправки

Heartbeat не контролирует отдельное сообщение напрямую, но позволяет определить потерю соединения.

Если heartbeat перестал приходить:

  • соединение считается неработоспособным;
  • клиент инициирует переподключение;
  • новые отправки блокируются или повторяются;
  • приложение может перевести сообщения в очередь ожидания.

Таким образом heartbeat создаёт базовый механизм таймаута для всей транспортной сессии.


Настройка агрессивных heartbeat-интервалов

Для систем реального времени используются небольшие интервалы:

const client = new Client({
    brokerURL: 'ws://localhost:15674/ws',

    heartbeatOutgoing: 3000,
    heartbeatIncoming: 3000
});

Преимущества:

  • быстрое обнаружение разрыва;
  • снижение времени потери сообщений;
  • оперативное переподключение.

Недостатки:

  • дополнительная нагрузка;
  • рост количества сетевых пакетов;
  • риск ложных срабатываний при медленной сети.

Реализация пользовательского таймаута отправки

Поскольку publish() не поддерживает Promise, таймаут обычно реализуется вручную.

Пример базового таймаута:

function publishWithTimeout(client, message, timeout = 5000) {
    return new Promise((resolve, reject) => {

        if (!client.connected) {
            reject(new Error('Нет подключения'));
            return;
        }

        const timer = setTimeout(() => {
            reject(new Error('Таймаут отправки'));
        }, timeout);

        try {
            client.publish(message);

            clearTimeout(timer);

            resolve();
        } catch (error) {
            clearTimeout(timer);

            reject(error);
        }
    });
}

Ограничения такого подхода

Этот механизм проверяет только успешность вызова publish().

Он не гарантирует:

  • получение сообщения брокером;
  • запись сообщения в очередь;
  • доставку подписчику;
  • отсутствие потери пакетов.

Поэтому для полноценного таймаута требуется подтверждение от сервера.


Таймауты через ACK

Надёжный способ контроля отправки — ожидание подтверждения.

Схема работы:

  1. клиент отправляет сообщение;
  2. сервер публикует ответ;
  3. клиент ожидает ACK-сообщение;
  4. если ACK не получен вовремя — срабатывает таймаут.

Пример подтверждаемой отправки

Отправка сообщения

const requestId = crypto.randomUUID();

client.publish({
    destination: '/app/send',
    headers: {
        'request-id': requestId
    },
    body: JSON.stringify({
        text: 'Привет'
    })
});

Подписка на подтверждения

const pendingRequests = new Map();

client.subscribe('/topic/ack', message => {

    const data = JSON.parse(message.body);

    const callback = pendingRequests.get(data.requestId);

    if (callback) {
        callback();

        pendingRequests.delete(data.requestId);
    }
});

Таймаут ожидания ACK

function sendWithAckTimeout(payload, timeout = 5000) {

    return new Promise((resolve, reject) => {

        const requestId = crypto.randomUUID();

        const timer = setTimeout(() => {

            pendingRequests.delete(requestId);

            reject(new Error('ACK timeout'));

        }, timeout);

        pendingRequests.set(requestId, () => {

            clearTimeout(timer);

            resolve();
        });

        client.publish({
            destination: '/app/send',
            headers: {
                'request-id': requestId
            },
            body: JSON.stringify(payload)
        });
    });
}

Контроль зависших сообщений

При большом количестве сообщений важно отслеживать зависшие операции.

Типичная структура:

const pendingMessages = new Map();

Каждое сообщение хранит:

  • идентификатор;
  • timestamp отправки;
  • timer;
  • retry count;
  • payload;
  • статус доставки.

Очистка просроченных сообщений

Пример периодической проверки:

setInterval(() => {

    const now = Date.now();

    for (const [id, message] of pendingMessages.entries()) {

        if (now - message.createdAt > 10000) {

            console.error('Сообщение просрочено:', id);

            pendingMessages.delete(id);
        }
    }

}, 1000);

Таймауты при использовании RabbitMQ

В связке STOMP.js и RabbitMQ возможны ситуации, когда:

  • WebSocket активен;
  • STOMP-сессия существует;
  • очередь перегружена;
  • сообщение не маршрутизируется;
  • exchange не найден;
  • queue отсутствует.

В таких условиях отправка может формально завершиться успешно, хотя сообщение будет потеряно.

Поэтому в RabbitMQ часто применяются:

  • publisher confirms;
  • dead-letter queues;
  • retry exchanges;
  • application ACK;
  • delayed retry.

Таймауты и reconnectDelay

Параметр reconnectDelay влияет на скорость восстановления после таймаута.

const client = new Client({
    brokerURL: 'ws://localhost:15674/ws',

    reconnectDelay: 5000
});

Поведение

После обнаружения проблемы:

  1. соединение закрывается;
  2. запускается ожидание;
  3. выполняется reconnect;
  4. восстанавливаются подписки;
  5. повторяются отправки.

Комбинация heartbeat и reconnectDelay

На практике обычно используются:

const client = new Client({

    brokerURL: 'ws://localhost:15674/ws',

    heartbeatIncoming: 4000,
    heartbeatOutgoing: 4000,

    reconnectDelay: 3000
});

Такая конфигурация:

  • обнаруживает проблему примерно за 4 секунды;
  • переподключается через 3 секунды;
  • минимизирует время простоя.

Таймауты повторной отправки

При временных сетевых ошибках используется retry-механизм.

Пример:

async function sendWithRetry(payload, retries = 3) {

    for (let i = 0; i < retries; i++) {

        try {

            await sendWithAckTimeout(payload);

            return;

        } catch (error) {

            console.error('Ошибка отправки:', error);

            if (i === retries - 1) {
                throw error;
            }
        }
    }
}

Экспоненциальная задержка повторов

Фиксированная задержка создаёт нагрузку на брокер. Более устойчивым считается exponential backoff.

function wait(ms) {
    return new Promise(resolve => setTimeout(resolve, ms));
}

async function sendWithBackoff(payload) {

    let delay = 1000;

    for (let i = 0; i < 5; i++) {

        try {

            await sendWithAckTimeout(payload);

            return;

        } catch (error) {

            await wait(delay);

            delay *= 2;
        }
    }

    throw new Error('Не удалось отправить сообщение');
}

Таймауты в offline-first приложениях

В offline-first архитектуре сообщения временно сохраняются локально.

Алгоритм:

  1. попытка отправки;
  2. ожидание ACK;
  3. таймаут;
  4. запись сообщения в локальное хранилище;
  5. повтор после reconnect.

Буферизация сообщений

Пример очереди:

const offlineQueue = [];

Добавление:

offlineQueue.push({
    body: payload,
    createdAt: Date.now()
});

Повторная отправка:

async function flushQueue() {

    while (offlineQueue.length > 0) {

        const item = offlineQueue.shift();

        await sendWithAckTimeout(item.body);
    }
}

Опасность бесконечных таймаутов

Отсутствие ограничений может привести к:

  • переполнению памяти;
  • накоплению таймеров;
  • тысячам pending Promise;
  • утечкам listeners;
  • перегрузке reconnect-попытками.

Поэтому всегда ограничиваются:

  • время ожидания;
  • количество retry;
  • размер offline-очереди;
  • максимальное число reconnect.

Таймауты и утечки памяти

Типичная ошибка:

setTimeout(() => {
    reject(new Error());
}, 5000);

Без очистки:

clearTimeout(timer);

таймер продолжит существовать даже после успешной отправки.

При тысячах сообщений это становится серьёзной проблемой.


Централизованный менеджер таймаутов

В крупных приложениях создаётся отдельный сервис управления отправкой.

Пример структуры:

class MessageTimeoutManager {

    constructor() {
        this.pending = new Map();
    }

    register(id, resolve, reject, timeout) {

        const timer = setTimeout(() => {

            this.pending.delete(id);

            reject(new Error('Timeout'));

        }, timeout);

        this.pending.set(id, {
            timer,
            resolve,
            reject
        });
    }

    complete(id) {

        const item = this.pending.get(id);

        if (!item) {
            return;
        }

        clearTimeout(item.timer);

        item.resolve();

        this.pending.delete(id);
    }
}

Таймауты в высоконагруженных системах

В системах с большим количеством сообщений учитываются:

  • latency сети;
  • скорость брокера;
  • размер payload;
  • количество подписчиков;
  • нагрузка на CPU;
  • backpressure;
  • ограничения WebSocket.

Неверно выбранный timeout способен:

  • вызывать ложные ошибки;
  • провоцировать повторные отправки;
  • создавать дубликаты;
  • перегружать брокер.

Рекомендуемые значения

Чат-приложения

2–5 секунд

Финансовые системы

1–3 секунды

IoT

10–30 секунд

Медленные мобильные сети

15–60 секунд

Диагностика таймаутов

Для анализа проблем обычно логируются:

  • timestamp отправки;
  • request id;
  • время ACK;
  • reconnect-события;
  • heartbeat timeout;
  • retry count;
  • длительность reconnect;
  • queue size.

Пример расширенного логирования

async function monitoredSend(payload) {

    const started = Date.now();

    try {

        await sendWithAckTimeout(payload);

        console.log('Отправлено за', Date.now() - started, 'ms');

    } catch (error) {

        console.error('Таймаут отправки', {
            duration: Date.now() - started,
            error
        });

        throw error;
    }
}

Типичная production-схема

В production-приложениях таймауты отправки обычно строятся как комбинация нескольких механизмов:

  1. heartbeat;
  2. reconnectDelay;
  3. ACK-подтверждения;
  4. retry;
  5. exponential backoff;
  6. offline queue;
  7. дедупликация сообщений;
  8. мониторинг зависших операций;
  9. централизованное управление таймерами;
  10. логирование latency и reconnect-событий.