Skip to content

OfflineMutationQueue

OfflineMutationQueue сохраняет разрешённые изменяющие операции в Web Storage и повторяет их после восстановления сети.

Очередь нужна только для команд, которые продукт разрешает выполнить offline. Обычный GET в очередь не помещается.

Что требуется от backend

Offline replay может повториться после reload, сбоя вкладки или сетевой ошибки. Backend endpoint должен поддерживать idempotency key и атомарно не выполнять одну команду дважды.

Локальная deduplication очереди не заменяет backend idempotency.

Тип переменных задачи

ts
export interface CreateCommentInput {
    postId: string
    text: string
}

export interface CreateCommentVariables {
    // Полезные данные команды.
    input: CreateCommentInput

    // Один key используется и online-запросом, и replay.
    idempotencyKey: string
}

// Проектный runtime guard для canPersistVariables.
export const isCreateCommentVariables = (value: unknown): value is CreateCommentVariables => {
    if (typeof value !== "object" || value === null) {
        return false
    }

    const record = value as Record<string, unknown>
    const input = record.input

    return (
        typeof record.idempotencyKey === "string" &&
        typeof input === "object" &&
        input !== null &&
        typeof (input as Record<string, unknown>).postId === "string" &&
        typeof (input as Record<string, unknown>).text === "string"
    )
}

Переменные должны быть plain JSON. File, Blob, FormData, Date, Map, Set, BigInt, functions и custom class instances запрещены.

API-функция

ts
export const createCommentApi = (
    variables: CreateCommentVariables,
    config?: AxiosRequestConfig,
): Promise<AxiosResponse<Comment>> => {
    return authApi.post<Comment>("/api/v1/comments", variables.input, {
        ...config,
        headers: {
            ...config?.headers,
            "Idempotency-Key": variables.idempotencyKey,
        },
    })
}

Создание очереди

ts
import { OfflineMutationQueue } from "@dubium/query-layer/offline"

import { runtime, appScope } from "./query-runtime"
import { createCommentApi } from "../data/comments/api"

export const offlineQueue = new OfflineMutationQueue(runtime, "PORTAL_OFFLINE_MUTATIONS", {
    // Runtime сам создаст hash этого partition для storage key.
    scope: {
        applicationId: "portal",
        tenantId: sessionStore.tenantId,
        userId: sessionStore.userId,
    },

    // Без whitelist enqueue запрещён для всех mutation keys.
    whitelist: ["comments.create"],

    // После пяти неуспешных replay задача уйдёт в dead letter.
    maxAttempts: 5,

    // Внутри одного replay HTTP можно повторить два раза.
    retry: 2,
    retryDelay: (failureCount) => failureCount * 1_000,

    storage: window.localStorage,

    // Дополнительная project policy payload.
    canPersistVariables: (mutationKey, variables) => {
        return mutationKey === "comments.create" && isCreateCommentVariables(variables)
    },
})

// Registry связывает persisted mutationKey с реальной API-функцией.
offlineQueue.register<CreateComment, CreateCommentVariables>("comments.create", async (variables, signal) => {
    const response = await createCommentApi(variables, { signal })

    // После replay инвалидируем logical key через scope API.
    await appScope.query.invalidate({ queryKey: ["comments", "list"] }, { refetchActive: true })

    return response
})

// mount восстанавливает storage, слушает online и сразу пробует replay.
export const unmountOfflineQueue = offlineQueue.mount()

Handler нужно зарегистрировать при каждом startup до replay. В persisted snapshot хранится mutationKey, но не функция.

Domain store: online и offline путь

ts
import { makeAutoObservable } from "mobx"
import { MutationStore } from "@dubium/query-layer"
import { OfflineMutationQueue } from "@dubium/query-layer/offline"

export class CreateCommentStore {
    private readonly requestHandler = new MutationStore<Comment, [CreateCommentVariables]>({
        mutationKey: "comments.create",
        request: (variables, config) => {
            return createCommentApi(variables, config)
        },
        options: {
            // Повтор команды контролирует idempotency contract приложения.
            retry: false,
        },
    })

    constructor(private readonly queue: OfflineMutationQueue) {
        makeAutoObservable<this, "queue" | "requestHandler">(
            this,
            {
                queue: false,
                requestHandler: false,
            },
            {
                autoBind: true,
            },
        )
    }

    get loading(): boolean {
        return this.requestHandler.loading
    }

    async create(input: CreateCommentInput): Promise<Comment | null> {
        const variables: CreateCommentVariables = {
            input,
            idempotencyKey: crypto.randomUUID(),
        }

        // Если browser уже знает, что сети нет, сразу сохраняем задачу.
        if (!navigator.onLine) {
            this.queue.enqueue("comments.create", variables, {
                idempotencyKey: variables.idempotencyKey,
            })
            return null
        }

        try {
            const response = await this.requestHandler.execute(variables)
            return response?.data ?? null
        } catch (error) {
            // navigator.onLine не гарантирует доступность backend.
            // После фактической network error сохраняем ту же команду с тем же key.
            if (this.requestHandler.error?.code === "ERR_NETWORK") {
                this.queue.enqueue("comments.create", variables, {
                    idempotencyKey: variables.idempotencyKey,
                })
                return null
            }

            // HTTP 400/403/422 — не offline. Отдаём ошибку выше.
            throw error
        }
    }
}

Не добавляйте несуществующий networkMode: "offlineFirst": публичный TNetworkMode содержит только online и always. Offline flow реализует очередь и domain store явно.

enqueue

ts
const task = offlineQueue.enqueue("comments.create", variables, {
    idempotencyKey: variables.idempotencyKey,
})

Перед сохранением очередь:

  1. проверяет whitelist;
  2. проверяет JSON-safe форму;
  3. отклоняет поля с именами password/token/secret и аналогичными;
  4. вызывает canPersistVariables, если он задан;
  5. проверяет duplicate idempotency key;
  6. клонирует variables и сохраняет snapshot.

Та же whitelist, expiration, JSON-safe и sensitive-data policy применяется при restore к активным и dead-letter задачам. retryDeadLetter() повторно проверяет эту policy, лимит очереди и duplicate idempotency key; при отказе запись остаётся в dead letter. Каждая HTTP-попытка replay проходит через общую concurrency-очередь того же QueryRuntime, поэтому учитывается maxConcurrentRequests.

Все options конструктора

OptionTypeDefaultЧто делает
scope{ applicationId; tenantId?; userId? }Создаёт partition hash storage key.
whitelistreadonly string[]все запрещеноРазрешённые mutation keys.
storageStoragelocalStorage в browserХранилище snapshot.
dedupebooleantrueОтклонять duplicate idempotency key.
getIdempotencyKey(mutationKey, variables) => stringСоздаёт key, если enqueue его не передал.
canPersistVariables(mutationKey, variables) => booleantrue после встроенного guardProject allowlist payload.
maxAgenumberбез лимитаУдалять слишком старые задачи.
maxAttemptsnumberбез лимитаПосле какого failureCount переносить в dead letter.
deadLetterLimitnumber100Максимальное число dead-letter записей.
queueLimitnumber1000Максимальное число активных задач в памяти/storage.
retryRetryValuefalseПовторы только для задач с idempotency key.
retryDelayRetryDelayretryer defaultЗадержка повторов.
lockManagerIQueryLockManagerНе даёт нескольким вкладкам replay одновременно.
replayLockNamestringpartitioned defaultИмя межвкладочной блокировки.
onConflict(task, error) => "drop" | "keep"keepРешает судьбу HTTP 409.
onDropTask(task, reason) => voidСообщает об expired/duplicate/sensitive/failed task.
onReplayError(task, error) => voidПолучает ошибку отдельного replay.
onPersistError(error) => voidПолучает storage/serialization error.
clockIClocksystem clockВремя для тестов.
idGeneratorIIdGeneratorruntime generatorTask ids.
timeoutManagerITimeoutManagersystem timersRetry timers.

Методы и observable state

ЧленЧто делает
queueObservable array активных задач.
deadLettersObservable array остановленных задач.
replayingИдёт ли replay сейчас.
initialize()Один раз восстанавливает snapshot из storage.
mount()Идемпотентно подписывается на online, запускает replay и возвращает один unmount-listener.
register(key, fn)Регистрирует обработчик persisted mutation key.
enqueue(key, variables, options)Добавляет JSON-safe задачу.
replay()Вручную запускает доступные задачи.
cancelReplay()Отменяет текущий replay, не удаляя задачи.
remove(taskId)Удаляет активную задачу.
getAll()Возвращает snapshot активных задач.
getDeadLetters()Возвращает snapshot dead letter.
retryDeadLetter(taskId)Возвращает задачу в активную очередь.
discardDeadLetter(taskId)Окончательно удаляет dead-letter задачу.

У класса нет методов unmount() и dispose(). Функцию снятия listener возвращает mount().

Lifecycle

ts
// Снимаем online listener.
unmountOfflineQueue()

// Останавливаем текущий replay. Задачи остаются в storage.
offlineQueue.cancelReplay()

// Затем освобождаем scope/runtime.
appScope.dispose()
await runtime.dispose()

На logout продукт должен явно решить: удалить задачи старого пользователя, завершить их до logout или оставить в partitioned storage для той же session.