diff --git a/packages/parser2/package.json b/packages/parser2/package.json index 04192c6..057fec9 100644 --- a/packages/parser2/package.json +++ b/packages/parser2/package.json @@ -1,6 +1,6 @@ { "name": "@coopenomics/parser2", - "version": "1.1.2", + "version": "1.1.3", "description": "Universal EOSIO/Antelope SHiP-to-Redis blockchain indexer (parser)", "license": "MIT", "author": "Coopenomics contributors", diff --git a/packages/parser2/src/adapters/IoRedisStore.ts b/packages/parser2/src/adapters/IoRedisStore.ts index c25a144..e1700ca 100644 --- a/packages/parser2/src/adapters/IoRedisStore.ts +++ b/packages/parser2/src/adapters/IoRedisStore.ts @@ -37,7 +37,9 @@ interface IRedisClient { count: string, countVal: number, block: string, blockMs: number, streams: string, stream: string, id: string, - ): Promise]> | null> + ): Promise]> | null> + // XPENDING key group (summary): [count, min-id, max-id, [[consumer, count],...]|null] + xpending(key: string, group: string): Promise<[number, string | null, string | null, Array<[string, string]> | null] | null> xrange(key: string, start: string, end: string, count: string, countVal: number): Promise> xrevrange(key: string, end: string, start: string, count: string, countVal: number): Promise> xlen(key: string): Promise @@ -94,9 +96,14 @@ return 0 * Redis возвращает: [[id, [field1, val1, field2, val2, ...]], ...] * Мы конвертируем в: [{id, fields: {field1: val1, ...}}, ...] */ -function parseStreamEntries(raw: Array<[string, string[]]>): StreamMessage[] { +function parseStreamEntries(raw: Array<[string, string[] | null]>): StreamMessage[] { const messages: StreamMessage[] = [] for (const [msgId, rawFields] of raw) { + // Обрезанные/удалённые записи: при чтении PEL (XREADGROUP ... 0) Redis + // отдаёт [id, null] для ID, которых уже нет в стриме (XTRIM/XDEL их снёс). + // rawFields === null → .length бросает TypeError и убивает consumer навсегда. + // Пропускаем — запись недоступна, восстановить нечего. + if (!rawFields) continue const fields: Record = {} // rawFields — плоский массив [key, val, key, val, ...], шагаем по 2 for (let i = 0; i + 1 < rawFields.length; i += 2) { @@ -246,6 +253,19 @@ export class IoRedisStore implements RedisStore { await this.client.xack(stream, group, id) } + /** + * XPENDING stream group (summary-форма) → самый старый un-acked ID в PEL группы. + * null если pending нет. Нужен для безопасного XTRIM: триммить можно только + * НИЖЕ этого ID, иначе снесём собственные недоставленные/неподтверждённые записи. + */ + async xpendingMinId(stream: string, group: string): Promise { + const res = await this.client.xpending(stream, group) + if (!Array.isArray(res)) return null + const [count, minId] = res + if (!count || !minId) return null + return String(minId) + } + /** ZADD key score member. */ async zadd(key: string, score: number, member: string): Promise { await this.client.zadd(key, score, member) diff --git a/packages/parser2/src/core/XtrimSupervisor.ts b/packages/parser2/src/core/XtrimSupervisor.ts index 4837d1c..9720888 100644 --- a/packages/parser2/src/core/XtrimSupervisor.ts +++ b/packages/parser2/src/core/XtrimSupervisor.ts @@ -7,9 +7,15 @@ * Стратегия MINID: * Вместо хранения фиксированного числа записей (MAXLEN), мы сохраняем все * записи, которые ещё не подтверждены (pending) хотя бы одной consumer group. - * minId = min(lastDeliveredId всех групп с pending > 0). + * minId = min(самый старый un-acked pending ID всех групп с pending > 0). * XTRIM stream MINID minId удаляет всё с ID < minId. * + * ВАЖНО: триммить по lastDeliveredId НЕЛЬЗЯ — un-acked pending записи старше + * lastDeliveredId (доставлены, но не подтверждены: их ID < last-delivered). + * XTRIM MINID lastDeliveredId снёс бы собственные pending группы → consumer при + * перечитывании PEL получает [id, null] и падает. Поэтому нижняя граница trim — + * самый старый un-acked ID из XPENDING, а не lastDeliveredId. + * * Это гарантирует, что ни один consumer не потеряет сообщения при trim: * группа с отставанием «тормозит» trim, пока не догонит. * @@ -18,6 +24,19 @@ import type { RedisStore } from '../ports/RedisStore.js' +/** + * Числовое сравнение Redis Stream ID формата `-`. + * Строковое сравнение (a < b) неверно: '1000-0' < '999-0' лексикографически + * true, хотя по времени 1000 > 999. Возвращает <0, 0, >0. + */ +function compareStreamIds(a: string, b: string): number { + const [aMs, aSeq = '0'] = a.split('-') + const [bMs, bSeq = '0'] = b.split('-') + const ms = Number(aMs) - Number(bMs) + if (ms !== 0) return ms + return Number(aSeq) - Number(bSeq) +} + export interface XtrimSupervisorOpts { redis: RedisStore /** Имя стрима для очистки (обычно ce:parser2::events). */ @@ -60,8 +79,8 @@ export class XtrimSupervisor { * Один цикл очистки: * 1. Получаем список consumer groups через XINFO GROUPS. * 2. Фильтруем группы у которых есть pending сообщения (pending > 0). - * 3. Находим минимальный lastDeliveredId среди таких групп. - * 4. XTRIM stream MINID minId — удаляем всё старее этого ID. + * 3. Для каждой берём самый старый un-acked pending ID (XPENDING). + * 4. minId = числовой минимум этих ID; XTRIM stream MINID minId. * * Если pending-групп нет — trim не делается (всё подтверждено). * Если стрим не существует или XInfo бросает — тихо игнорируем (best-effort). @@ -75,12 +94,19 @@ export class XtrimSupervisor { const pendingGroups = groups.filter(g => g.pending > 0) if (pendingGroups.length === 0) return - // Наименьший lastDeliveredId = самый отстающий consumer - const minId = pendingGroups - .map(g => g.lastDeliveredId) - .reduce((a, b) => (a < b ? a : b)) + // Самый старый un-acked pending ID по каждой группе. НЕ lastDeliveredId: + // он выше pending, trim по нему снёс бы собственные неподтверждённые записи. + const oldestPending: string[] = [] + for (const g of pendingGroups) { + const id = await this.redis.xpendingMinId(this.stream, g.name) + if (id) oldestPending.push(id) + } + if (oldestPending.length === 0) return - if (minId) await this.redis.xtrim(this.stream, minId) + // Нижняя граница trim = самый старый un-acked ID среди всех групп + const minId = oldestPending.reduce((a, b) => (compareStreamIds(a, b) <= 0 ? a : b)) + + await this.redis.xtrim(this.stream, minId) } catch { // XTRIM — best-effort: ошибки не должны влиять на основной поток обработки } diff --git a/packages/parser2/src/ports/RedisStore.ts b/packages/parser2/src/ports/RedisStore.ts index a6fd67f..3f6cb88 100644 --- a/packages/parser2/src/ports/RedisStore.ts +++ b/packages/parser2/src/ports/RedisStore.ts @@ -79,6 +79,13 @@ export interface RedisStore { /** XACK stream group id — подтверждает обработку, убирает из PEL. */ xack(stream: string, group: string, id: string): Promise + /** + * XPENDING stream group → ID самого старого un-acked сообщения в PEL группы + * (или null если pending нет). Используется XtrimSupervisor: trim допустим + * только НИЖЕ этого ID, иначе un-acked записи будут потеряны. + */ + xpendingMinId(stream: string, group: string): Promise + // ── Sorted Set ──────────────────────────────────────────────────────────── /** ZADD key score member. */ diff --git a/packages/parser2/test/integration/xtrim.integration.test.ts b/packages/parser2/test/integration/xtrim.integration.test.ts new file mode 100644 index 0000000..a113093 --- /dev/null +++ b/packages/parser2/test/integration/xtrim.integration.test.ts @@ -0,0 +1,103 @@ +/** + * Integration тест против реального Redis: регресс-гард на баг #3. + * + * Баг: XtrimSupervisor триммил по lastDeliveredId и сносил собственные un-acked + * pending записи; при перечитывании PEL Redis отдавал [id, null], а + * parseStreamEntries падал на null.length → consumer умирал навсегда. + * + * Здесь воспроизводим точный путь на живом Redis (chain/SHiP не нужны): + * A. trim по lastDeliveredId зануляет pending → xreadGroup(PEL) НЕ падает на null. + * B. XtrimSupervisor.trim() (новая логика) сохраняет все un-acked pending. + * + * Требование: Redis на REDIS_URL (в CI всегда есть). Без Redis — тесты skip. + */ + +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { IoRedisStore } from '../../src/adapters/IoRedisStore.js' +import { XtrimSupervisor } from '../../src/core/XtrimSupervisor.js' + +const REDIS_URL = process.env['REDIS_URL'] ?? 'redis://localhost:6379/14' +const STREAM = '__it_bug3__:stream' +const GROUP = 'g' + +let store: IoRedisStore +let redisUp = false + +// XtrimSupervisor.trim приватный — дёргаем один цикл напрямую через каст. +function runTrimOnce(sup: XtrimSupervisor): Promise { + return (sup as unknown as { trim(): Promise }).trim() +} + +beforeAll(async () => { + store = new IoRedisStore({ url: REDIS_URL }) + try { + await store.connect() + await store.client.del(STREAM) + redisUp = true + } catch { + redisUp = false + } +}) + +afterAll(async () => { + if (redisUp) { + await store.client.del(STREAM) + await store.quit() + } +}) + +describe('bug #3 — XTRIM/PEL на реальном Redis', () => { + it('xreadGroup(PEL) не падает на обрезанных [id,null] записях', async () => { + if (!redisUp) return + await store.client.del(STREAM) + await store.xgroupCreate(STREAM, GROUP, '0') + for (let i = 0; i < 10; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) }) + + // читаем 5 → они pending (un-acked); lastDeliveredId = id 5-й записи + const read = await store.xreadGroup(STREAM, GROUP, 'c1', 5, 0, '>') + expect(read).toHaveLength(5) + const lastDelivered = read[4]!.id + + // ИМИТАЦИЯ СТАРОГО БАГА: trim по lastDeliveredId сносит первые 4 pending + await store.client.xtrim(STREAM, 'MINID', lastDelivered) + + // consumer перечитывает PEL — Redis отдаёт [id,null] для снесённых. + // Не должно бросать (старый parseStreamEntries падал на null.length). + const pel = await store.xreadGroup(STREAM, GROUP, 'c1', 100, 0, '0') + // обрезанные пропущены, выжившие — с полями + expect(pel.every(m => Object.keys(m.fields).length > 0)).toBe(true) + expect(pel.some(m => m.id === lastDelivered)).toBe(true) + }) + + it('XtrimSupervisor.trim() сохраняет все un-acked pending записи', async () => { + if (!redisUp) return + await store.client.del(STREAM) + await store.xgroupCreate(STREAM, GROUP, '0') + for (let i = 0; i < 10; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) }) + + const read = await store.xreadGroup(STREAM, GROUP, 'c1', 5, 0, '>') + expect(read).toHaveLength(5) + const oldestPending = read[0]!.id + + const sup = new XtrimSupervisor({ redis: store, stream: STREAM, intervalMs: 999_999 }) + await runTrimOnce(sup) + + // consumer перечитывает PEL — все 5 pending целы, без null, без краша + const pel = await store.xreadGroup(STREAM, GROUP, 'c1', 100, 0, '0') + expect(pel).toHaveLength(5) + expect(pel.every(m => Object.keys(m.fields).length > 0)).toBe(true) + // нижняя граница trim = oldestPending → первый pending на месте + expect(pel[0]!.id).toBe(oldestPending) + }) + + it('xpendingMinId возвращает самый старый un-acked ID, null без pending', async () => { + if (!redisUp) return + await store.client.del(STREAM) + await store.xgroupCreate(STREAM, GROUP, '0') + expect(await store.xpendingMinId(STREAM, GROUP)).toBeNull() + + for (let i = 0; i < 3; i++) await store.xadd(STREAM, { kind: 'action', n: String(i) }) + const read = await store.xreadGroup(STREAM, GROUP, 'c1', 3, 0, '>') + expect(await store.xpendingMinId(STREAM, GROUP)).toBe(read[0]!.id) + }) +}) diff --git a/packages/parser2/test/unit/ioRedisStore.test.ts b/packages/parser2/test/unit/ioRedisStore.test.ts index 2dd4989..2f093e7 100644 --- a/packages/parser2/test/unit/ioRedisStore.test.ts +++ b/packages/parser2/test/unit/ioRedisStore.test.ts @@ -5,6 +5,8 @@ vi.mock('ioredis', () => { const mockRedis = { xadd: vi.fn().mockResolvedValue('1700000000000-0'), xtrim: vi.fn().mockResolvedValue(5), + xreadgroup: vi.fn().mockResolvedValue(null), + xpending: vi.fn().mockResolvedValue([0, null, null, null]), zadd: vi.fn().mockResolvedValue(1), zrangebyscore: vi.fn().mockResolvedValue(['{"version":"eosio::abi/1.0"}']), zrevrangebyscore: vi.fn().mockResolvedValue(['{"version":"eosio::abi/1.0"}']), @@ -37,6 +39,34 @@ describe('IoRedisStore', () => { expect(store.client.xtrim).toHaveBeenCalledWith('mystream', 'MINID', '1234567890000-0') }) + it('xreadGroup skips trimmed entries with null fields (no crash — баг #3)', async () => { + // Redis отдаёт [id, null] для обрезанных/удалённых ID при чтении PEL. + // Раньше rawFields.length на null убивал consumer навсегда. + vi.mocked(store.client.xreadgroup).mockResolvedValueOnce([ + ['mystream', [ + ['100-0', ['kind', 'action']], + ['101-0', null], // обрезанная запись (XTRIM/XDEL) + ['102-0', ['kind', 'delta']], + ]], + ]) + const msgs = await store.xreadGroup('mystream', 'g', 'c', 10, 0, '0') + expect(msgs).toHaveLength(2) + expect(msgs.map(m => m.id)).toEqual(['100-0', '102-0']) + expect(msgs[0]!.fields).toEqual({ kind: 'action' }) + }) + + it('xpendingMinId returns oldest pending id from XPENDING summary', async () => { + vi.mocked(store.client.xpending).mockResolvedValueOnce([5, '90-0', '100-0', [['c', '5']]]) + const id = await store.xpendingMinId('mystream', 'g') + expect(id).toBe('90-0') + }) + + it('xpendingMinId returns null when group has no pending', async () => { + vi.mocked(store.client.xpending).mockResolvedValueOnce([0, null, null, null]) + const id = await store.xpendingMinId('mystream', 'g') + expect(id).toBeNull() + }) + it('zadd calls redis.zadd with score and member', async () => { await store.zadd('parser2:abi:eosio', 100, '{"version":"eosio::abi/1.0"}') expect(store.client.zadd).toHaveBeenCalledWith('parser2:abi:eosio', 100, '{"version":"eosio::abi/1.0"}') diff --git a/packages/parser2/test/unit/parserClient.test.ts b/packages/parser2/test/unit/parserClient.test.ts index 03d5ab4..f7c4e72 100644 --- a/packages/parser2/test/unit/parserClient.test.ts +++ b/packages/parser2/test/unit/parserClient.test.ts @@ -69,6 +69,7 @@ class FakeRedis implements RedisStore { xlen = vi.fn(async (): Promise => 0) xdel = vi.fn(async (): Promise => 0) xack = vi.fn(async (): Promise => {}) + xpendingMinId = vi.fn(async (): Promise => null) zadd = vi.fn(async (key: string, score: number, member: string): Promise => { if (!this.zsets.has(key)) this.zsets.set(key, new Map()) diff --git a/packages/parser2/test/unit/xtrimSupervisor.test.ts b/packages/parser2/test/unit/xtrimSupervisor.test.ts index c845493..fbfd2b6 100644 --- a/packages/parser2/test/unit/xtrimSupervisor.test.ts +++ b/packages/parser2/test/unit/xtrimSupervisor.test.ts @@ -3,20 +3,28 @@ * * Проверяем: * - start/stop идемпотентны - * - trim использует MINID по отстающей группе с pending > 0 + * - trim использует MINID по самому старому un-acked pending ID (XPENDING), + * а НЕ по lastDeliveredId (иначе сносит собственные pending — баг #3) + * - числовое сравнение ID (а не лексикографическое) * - группы без pending не мешают trim'у - * - XTRIM не вызывается если групп нет или ни у одной нет pending - * - ошибки в xinfoGroups не прерывают supervisor + * - XTRIM не вызывается если групп нет, ни у одной нет pending, или pending-ID null + * - ошибки в xinfoGroups/xtrim не прерывают supervisor */ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest' import { XtrimSupervisor } from '../../src/core/XtrimSupervisor.js' import type { RedisStore, XGroupInfo } from '../../src/ports/RedisStore.js' -function makeRedis(groups: XGroupInfo[] = []): RedisStore { +function makeRedis( + groups: XGroupInfo[] = [], + pendingIds: Record = {}, +): RedisStore { return { xinfoGroups: vi.fn().mockResolvedValue(groups), xtrim: vi.fn().mockResolvedValue(0), + xpendingMinId: vi.fn().mockImplementation((_stream: string, group: string) => + Promise.resolve(pendingIds[group] ?? null), + ), // остальные методы не используются в supervisor } as unknown as RedisStore } @@ -25,6 +33,7 @@ function makeRedis(groups: XGroupInfo[] = []): RedisStore { async function flushPromises(): Promise { await Promise.resolve() await Promise.resolve() + await Promise.resolve() } describe('XtrimSupervisor — lifecycle', () => { @@ -124,36 +133,95 @@ describe('XtrimSupervisor — trim logic', () => { sup.stop() }) - it('trims to minimum lastDeliveredId among groups with pending > 0', async () => { - // g1 pending но отстал на 100 — должен стать minId - // g2 pending и догнал до 500 — игнорируется для выбора min - // g3 нет pending — вообще не учитывается - const redis = makeRedis([ - { name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 }, - { name: 'g2', pending: 2, lastDeliveredId: '500-0', lag: 2, consumers: 1 }, - { name: 'g3', pending: 0, lastDeliveredId: '50-0', lag: 0, consumers: 1 }, - ]) + it('trims to the OLDEST un-acked pending ID (XPENDING), NOT lastDeliveredId', async () => { + // Регресс-гард на баг #3: pending-записи старше lastDeliveredId. + // g1: lastDelivered=100-0, но самый старый un-acked pending = 90-0. + // Триммить надо по 90-0, иначе un-acked 90..99 теряются и consumer падает. + const redis = makeRedis( + [{ name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 }], + { g1: '90-0' }, + ) const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 }) sup.start() await vi.advanceTimersByTimeAsync(100) await flushPromises() expect(redis.xtrim).toHaveBeenCalledTimes(1) - expect(redis.xtrim).toHaveBeenCalledWith('s', '100-0') + expect(redis.xtrim).toHaveBeenCalledWith('s', '90-0') + // именно НЕ lastDeliveredId + expect(redis.xtrim).not.toHaveBeenCalledWith('s', '100-0') sup.stop() }) - it('uses lexicographic comparison for stream IDs (which is also numeric-correct for same-length)', async () => { - const redis = makeRedis([ - { name: 'g1', pending: 1, lastDeliveredId: '1000-0', lag: 1, consumers: 1 }, - { name: 'g2', pending: 1, lastDeliveredId: '999-0', lag: 1, consumers: 1 }, - ]) + it('trims to minimum oldest-pending-id among groups with pending > 0', async () => { + // g1 отстал (oldest pending 90-0) → станет minId + // g2 догнал (oldest pending 480-0) → не минимум + // g3 нет pending → не учитывается (xpendingMinId не зовётся) + const redis = makeRedis( + [ + { name: 'g1', pending: 5, lastDeliveredId: '100-0', lag: 5, consumers: 1 }, + { name: 'g2', pending: 2, lastDeliveredId: '500-0', lag: 2, consumers: 1 }, + { name: 'g3', pending: 0, lastDeliveredId: '50-0', lag: 0, consumers: 1 }, + ], + { g1: '90-0', g2: '480-0' }, + ) const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 }) sup.start() await vi.advanceTimersByTimeAsync(100) await flushPromises() - // lexicographic: '1000-0' < '999-0' → выбран '1000-0' - expect(redis.xtrim).toHaveBeenCalledWith('s', '1000-0') + + expect(redis.xtrim).toHaveBeenCalledTimes(1) + expect(redis.xtrim).toHaveBeenCalledWith('s', '90-0') + // g3 без pending — XPENDING не запрашивался + expect(redis.xpendingMinId).not.toHaveBeenCalledWith('s', 'g3') + sup.stop() + }) + + it('uses NUMERIC (not lexicographic) comparison for stream IDs', async () => { + // Лексикографически '1000-0' < '999-0' (неверно). Численно 999 < 1000. + const redis = makeRedis( + [ + { name: 'g1', pending: 1, lastDeliveredId: '1000-0', lag: 1, consumers: 1 }, + { name: 'g2', pending: 1, lastDeliveredId: '999-0', lag: 1, consumers: 1 }, + ], + { g1: '1000-0', g2: '999-0' }, + ) + const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 }) + sup.start() + await vi.advanceTimersByTimeAsync(100) + await flushPromises() + // численный минимум = '999-0' + expect(redis.xtrim).toHaveBeenCalledWith('s', '999-0') + sup.stop() + }) + + it('compares by sequence when ms part is equal', async () => { + const redis = makeRedis( + [ + { name: 'g1', pending: 1, lastDeliveredId: '5-9', lag: 1, consumers: 1 }, + { name: 'g2', pending: 1, lastDeliveredId: '5-2', lag: 1, consumers: 1 }, + ], + { g1: '5-9', g2: '5-2' }, + ) + const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 }) + sup.start() + await vi.advanceTimersByTimeAsync(100) + await flushPromises() + expect(redis.xtrim).toHaveBeenCalledWith('s', '5-2') + sup.stop() + }) + + it('skips trim when pending groups exist but XPENDING returns null (race)', async () => { + // pending>0 в XINFO, но к моменту XPENDING всё подтверждено → null → не триммим + const redis = makeRedis( + [{ name: 'g1', pending: 3, lastDeliveredId: '100-0', lag: 3, consumers: 1 }], + { g1: null }, + ) + const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 }) + sup.start() + await vi.advanceTimersByTimeAsync(100) + await flushPromises() + expect(redis.xtrim).not.toHaveBeenCalled() sup.stop() }) @@ -173,9 +241,10 @@ describe('XtrimSupervisor — trim logic', () => { }) it('silently swallows errors from xtrim itself', async () => { - const redis = makeRedis([ - { name: 'g1', pending: 1, lastDeliveredId: '100-0', lag: 1, consumers: 1 }, - ]) + const redis = makeRedis( + [{ name: 'g1', pending: 1, lastDeliveredId: '100-0', lag: 1, consumers: 1 }], + { g1: '95-0' }, + ) ;(redis.xtrim as ReturnType).mockRejectedValue(new Error('xtrim failed')) const sup = new XtrimSupervisor({ redis, stream: 's', intervalMs: 100 })