Merge pull request 'feat(backpressure): окно поглощения writer↔consumer (защита Redis RAM на репарсинге)' (#10) from dev into main
Release / Release (push) Successful in 9m33s
Release / Release (push) Successful in 9m33s
This commit was merged in pull request #10.
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@coopenomics/parser2",
|
||||
"version": "1.1.3",
|
||||
"version": "1.2.0",
|
||||
"description": "Universal EOSIO/Antelope SHiP-to-Redis blockchain indexer (parser)",
|
||||
"license": "MIT",
|
||||
"author": "Coopenomics contributors",
|
||||
|
||||
@@ -43,6 +43,7 @@ interface IRedisClient {
|
||||
xrange(key: string, start: string, end: string, count: string, countVal: number): Promise<Array<[string, string[]]>>
|
||||
xrevrange(key: string, end: string, start: string, count: string, countVal: number): Promise<Array<[string, string[]]>>
|
||||
xlen(key: string): Promise<number>
|
||||
del(...keys: string[]): Promise<number>
|
||||
xdel(key: string, ...ids: string[]): Promise<number>
|
||||
xack(stream: string, group: string, id: string): Promise<number>
|
||||
zadd(key: string, score: number, member: string): Promise<number>
|
||||
|
||||
@@ -40,6 +40,12 @@ export interface ParserOptions {
|
||||
abiFallback?: 'rpc-current' | 'fail'
|
||||
/** XtrimSupervisor: интервал проверки и включение/отключение автообрезки стрима. */
|
||||
xtrim?: { intervalMs?: number; enabled?: boolean }
|
||||
/**
|
||||
* Backpressure-окно: пауза чтения SHiP когда backlog стрима (XLEN) достигает
|
||||
* highWater, возобновление при сливе до lowWater. Защита Redis RAM на
|
||||
* репарсинге, когда consumer медленнее писателя. По умолчанию включено.
|
||||
*/
|
||||
backpressure?: { enabled?: boolean; highWater?: number; lowWater?: number; pollMs?: number }
|
||||
/** ReconnectSupervisor: максимум попыток и backoff-таблица в секундах. */
|
||||
reconnect?: { maxAttempts?: number; backoffSeconds?: number[] }
|
||||
/** Десериализатор ABI-данных. Единственный вариант — 'wharfkit'. */
|
||||
|
||||
@@ -82,6 +82,26 @@ export const configSchema = {
|
||||
},
|
||||
additionalProperties: false,
|
||||
},
|
||||
/**
|
||||
* Backpressure: окно поглощения между писателем (SHiP) и consumer'ом.
|
||||
* Когда backlog стрима (XLEN) ≥ highWater — чтение блоков встаёт на паузу,
|
||||
* пока consumer не сольёт backlog ≤ lowWater. Бережёт Redis RAM на длинном
|
||||
* репарсинге. Ничего не теряет — события остаются в стриме.
|
||||
*/
|
||||
backpressure: {
|
||||
type: 'object',
|
||||
properties: {
|
||||
/** Включить/отключить. По умолчанию true. */
|
||||
enabled: { type: 'boolean', default: true },
|
||||
/** Верхняя граница backlog (XLEN) — пауза чтения. По умолчанию 100000. */
|
||||
highWater: { type: 'number', default: 100000 },
|
||||
/** Нижняя граница — возобновление. Должна быть < highWater. По умолчанию 50000. */
|
||||
lowWater: { type: 'number', default: 50000 },
|
||||
/** Интервал опроса XLEN во время паузы, мс. По умолчанию 200. */
|
||||
pollMs: { type: 'number', default: 200 },
|
||||
},
|
||||
additionalProperties: false,
|
||||
},
|
||||
/** ReconnectSupervisor: поведение при разрыве SHiP соединения. */
|
||||
reconnect: {
|
||||
type: 'object',
|
||||
|
||||
@@ -0,0 +1,100 @@
|
||||
/**
|
||||
* Backpressure-окно между чтением блокчейна (писатель) и consumer'ом (читатель).
|
||||
*
|
||||
* Проблема: Parser читает блоки из SHiP и XADD'ит события в Redis-стрим быстрее,
|
||||
* чем ParserClient (controller) их вычитывает. На длинном репарсинге (genesis →
|
||||
* head, миллионы блоков) стрим растёт в ПАМЯТИ REDIS неограниченно → Redis OOM.
|
||||
* RAM самого парсера ограничена SHiP-окном (max_messages_in_flight), а вот
|
||||
* Redis-стрим — нет. XtrimSupervisor не спасает: он не режет непрочитанное.
|
||||
*
|
||||
* Решение: перед обработкой каждого блока проверяем backlog = XLEN стрима.
|
||||
* Если backlog >= highWater — СТАВИМ чтение на паузу (не ack'аем блок, SHiP
|
||||
* перестаёт слать) и ждём пока consumer сольёт backlog до lowWater. Потом
|
||||
* продолжаем. Ничего не теряем: события лежат в стриме, мы лишь не доливаем.
|
||||
*
|
||||
* Если consumer вообще не подключён — backlog растёт до highWater и писатель
|
||||
* замирает, держа Redis в пределах ~highWater событий.
|
||||
*
|
||||
* Почему XLEN, а не lag группы: XLEN — это ровно тот объём, что занят в Redis
|
||||
* RAM (единственный защищаемый ресурс). Не зависит от числа/наличия групп и
|
||||
* версии Redis (lag появился в 7.0). Консервативен: считает и непрочитанное,
|
||||
* и ещё не обрезанное XtrimSupervisor'ом — пауза чуть раньше, это безопасно.
|
||||
*/
|
||||
|
||||
import type { RedisStore } from '../ports/RedisStore.js'
|
||||
import { rootLogger, type Logger } from '../logger.js'
|
||||
|
||||
const sleep = (ms: number): Promise<void> => new Promise(resolve => setTimeout(resolve, ms))
|
||||
|
||||
export interface BackpressureGateOpts {
|
||||
redis: RedisStore
|
||||
/** Стрим, backlog которого ограничиваем (ce:parser2:<chainId>:events). */
|
||||
stream: string
|
||||
/** Достигли этого XLEN — пауза чтения SHiP. */
|
||||
highWater: number
|
||||
/** Возобновляем чтение когда backlog слился до этого XLEN. < highWater. */
|
||||
lowWater: number
|
||||
/** Интервал опроса XLEN во время паузы, мс. По умолчанию 200. */
|
||||
pollMs?: number
|
||||
/** Логгер (по умолчанию rootLogger). */
|
||||
logger?: Logger
|
||||
}
|
||||
|
||||
export class BackpressureGate {
|
||||
private readonly redis: RedisStore
|
||||
private readonly stream: string
|
||||
private readonly highWater: number
|
||||
private readonly lowWater: number
|
||||
private readonly pollMs: number
|
||||
private readonly log: Logger
|
||||
private pausedNow = false
|
||||
|
||||
constructor(opts: BackpressureGateOpts) {
|
||||
if (!Number.isFinite(opts.highWater) || opts.highWater <= 0) {
|
||||
throw new Error(`BackpressureGate: highWater must be > 0, got ${opts.highWater}`)
|
||||
}
|
||||
if (!Number.isFinite(opts.lowWater) || opts.lowWater < 0 || opts.lowWater >= opts.highWater) {
|
||||
throw new Error(`BackpressureGate: lowWater must be in [0, highWater), got ${opts.lowWater} (highWater=${opts.highWater})`)
|
||||
}
|
||||
this.redis = opts.redis
|
||||
this.stream = opts.stream
|
||||
this.highWater = opts.highWater
|
||||
this.lowWater = opts.lowWater
|
||||
this.pollMs = opts.pollMs ?? 200
|
||||
this.log = (opts.logger ?? rootLogger).child({ component: 'BackpressureGate', stream: opts.stream })
|
||||
}
|
||||
|
||||
/** true пока чтение блоков стоит на паузе (для метрик/health). */
|
||||
get isPaused(): boolean {
|
||||
return this.pausedNow
|
||||
}
|
||||
|
||||
/**
|
||||
* Блокирует выполнение пока backlog (XLEN стрима) >= highWater.
|
||||
* Возвращается когда backlog опустился < lowWater ИЛИ shouldStop() стал true
|
||||
* (graceful shutdown). Если backlog уже < highWater — возвращается сразу,
|
||||
* не делая лишних RTT после первого XLEN.
|
||||
*
|
||||
* @param shouldStop предикат прерывания (обычно () => this.stopSignal)
|
||||
*/
|
||||
async waitForCapacity(shouldStop: () => boolean = () => false): Promise<void> {
|
||||
let backlog = await this.redis.xlen(this.stream)
|
||||
if (backlog < this.highWater) return
|
||||
|
||||
this.pausedNow = true
|
||||
this.log.warn(
|
||||
{ backlog, highWater: this.highWater, lowWater: this.lowWater },
|
||||
'backpressure: чтение блоков на паузе — consumer отстаёт, ждём слива стрима',
|
||||
)
|
||||
try {
|
||||
while (!shouldStop()) {
|
||||
await sleep(this.pollMs)
|
||||
backlog = await this.redis.xlen(this.stream)
|
||||
if (backlog <= this.lowWater) break
|
||||
}
|
||||
} finally {
|
||||
this.pausedNow = false
|
||||
}
|
||||
this.log.info({ backlog, lowWater: this.lowWater }, 'backpressure: backlog слит, чтение возобновлено')
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ import { WorkerPool } from '../workers/WorkerPool.js'
|
||||
import { BlockProcessor } from './BlockProcessor.js'
|
||||
import type { XtrimSupervisorOpts } from './XtrimSupervisor.js'
|
||||
import { XtrimSupervisor } from './XtrimSupervisor.js'
|
||||
import { BackpressureGate } from './BackpressureGate.js'
|
||||
import { ForkDetector } from './ForkDetector.js'
|
||||
import { RedisKeys } from '../redis/keys.js'
|
||||
import { ChainIdMismatchError } from '../errors.js'
|
||||
@@ -20,6 +21,7 @@ export class Parser {
|
||||
private workerPool: WorkerPool | null = null
|
||||
private blockProcessor: BlockProcessor | null = null
|
||||
private xtrimSupervisor: XtrimSupervisor | null = null
|
||||
private backpressureGate: BackpressureGate | null = null
|
||||
private running = false
|
||||
private stopSignal = false
|
||||
|
||||
@@ -100,6 +102,19 @@ export class Parser {
|
||||
this.xtrimSupervisor.start()
|
||||
}
|
||||
|
||||
// Backpressure-окно: пауза чтения SHiP когда backlog стрима упирается в
|
||||
// highWater, чтобы Redis не рос неограниченно на репарсинге. По умолчанию вкл.
|
||||
if (this.opts.backpressure?.enabled !== false) {
|
||||
const bp = this.opts.backpressure
|
||||
this.backpressureGate = new BackpressureGate({
|
||||
redis: this.redis,
|
||||
stream: eventsStream,
|
||||
highWater: bp?.highWater ?? 100_000,
|
||||
lowWater: bp?.lowWater ?? 50_000,
|
||||
...(bp?.pollMs !== undefined ? { pollMs: bp.pollMs } : {}),
|
||||
})
|
||||
}
|
||||
|
||||
const streamOpts = {
|
||||
startBlock: havePositions[0]?.blockNum ?? 0,
|
||||
havePositions,
|
||||
@@ -111,6 +126,14 @@ export class Parser {
|
||||
for await (const block of this.chainClient.streamBlocks(streamOpts)) {
|
||||
if (this.stopSignal) break
|
||||
|
||||
// Backpressure: пока backlog стрима выше окна — ждём, не доливаем и не
|
||||
// ack'аем (SHiP сам притормозит после max_messages_in_flight). Блок уже в
|
||||
// RAM (≤ окна SHiP), обработаем его как только consumer освободит место.
|
||||
if (this.backpressureGate) {
|
||||
await this.backpressureGate.waitForCapacity(() => this.stopSignal)
|
||||
if (this.stopSignal) break
|
||||
}
|
||||
|
||||
if (irreversibleOnly && block.thisBlock.blockNum > block.lastIrreversible.blockNum) {
|
||||
this.chainClient.ack(1)
|
||||
continue
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
/**
|
||||
* Integration тест против реального Redis: backpressure-окно.
|
||||
*
|
||||
* Проверяем на живом Redis (chain/SHiP не нужны):
|
||||
* - backlog < highWater → waitForCapacity не паузит, возвращается сразу
|
||||
* - backlog >= highWater → пауза; после слива стрима до <= lowWater — возврат
|
||||
*
|
||||
* Требование: Redis на REDIS_URL (в CI всегда есть). Без Redis — тесты skip.
|
||||
*/
|
||||
|
||||
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
|
||||
import { IoRedisStore } from '../../src/adapters/IoRedisStore.js'
|
||||
import { BackpressureGate } from '../../src/core/BackpressureGate.js'
|
||||
import type { Logger } from '../../src/logger.js'
|
||||
|
||||
const REDIS_URL = process.env['REDIS_URL'] ?? 'redis://localhost:6379/14'
|
||||
const STREAM = '__it_bp__:stream'
|
||||
|
||||
const silentLog = {
|
||||
child: () => silentLog, warn: () => {}, info: () => {}, debug: () => {}, error: () => {},
|
||||
} as unknown as Logger
|
||||
|
||||
const sleep = (ms: number): Promise<void> => new Promise(r => setTimeout(r, ms))
|
||||
|
||||
let store: IoRedisStore
|
||||
let redisUp = false
|
||||
|
||||
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('backpressure — окно на реальном Redis', () => {
|
||||
it('не паузит когда backlog < highWater', async () => {
|
||||
if (!redisUp) return
|
||||
await store.client.del(STREAM)
|
||||
await store.xadd(STREAM, { n: '1' })
|
||||
const gate = new BackpressureGate({ redis: store, stream: STREAM, highWater: 5, lowWater: 2, pollMs: 10, logger: silentLog })
|
||||
await gate.waitForCapacity()
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
|
||||
it('паузит при backlog >= highWater и возобновляет после слива до <= lowWater', async () => {
|
||||
if (!redisUp) return
|
||||
await store.client.del(STREAM)
|
||||
const ids: string[] = []
|
||||
for (let i = 0; i < 6; i++) ids.push(await store.xadd(STREAM, { n: String(i) }))
|
||||
expect(await store.xlen(STREAM)).toBe(6)
|
||||
|
||||
const gate = new BackpressureGate({ redis: store, stream: STREAM, highWater: 5, lowWater: 2, pollMs: 10, logger: silentLog })
|
||||
let resolved = false
|
||||
const p = gate.waitForCapacity().then(() => { resolved = true })
|
||||
|
||||
// даём gate войти в паузу и подождать
|
||||
await sleep(40)
|
||||
expect(gate.isPaused).toBe(true)
|
||||
expect(resolved).toBe(false)
|
||||
|
||||
// сливаем backlog: удаляем 4 записи → XLEN=2 <= lowWater
|
||||
await store.client.xdel(STREAM, ids[0]!, ids[1]!, ids[2]!, ids[3]!)
|
||||
expect(await store.xlen(STREAM)).toBe(2)
|
||||
|
||||
await p // должно разблокироваться
|
||||
expect(resolved).toBe(true)
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,126 @@
|
||||
/**
|
||||
* Тесты BackpressureGate — окна поглощения между писателем и consumer'ом.
|
||||
*
|
||||
* Проверяем:
|
||||
* - валидацию highWater/lowWater в конструкторе
|
||||
* - мгновенный возврат когда backlog < highWater (без паузы)
|
||||
* - паузу при backlog >= highWater и возобновление при сливе до <= lowWater
|
||||
* - isPaused отражает состояние
|
||||
* - shouldStop() прерывает паузу (graceful shutdown) даже если backlog высок
|
||||
*/
|
||||
|
||||
import { describe, it, expect, vi } from 'vitest'
|
||||
import { BackpressureGate } from '../../src/core/BackpressureGate.js'
|
||||
import type { RedisStore } from '../../src/ports/RedisStore.js'
|
||||
import type { Logger } from '../../src/logger.js'
|
||||
|
||||
// Молчаливый логгер — не засоряет вывод тестов.
|
||||
const silentLog = {
|
||||
child: () => silentLog,
|
||||
warn: () => {},
|
||||
info: () => {},
|
||||
debug: () => {},
|
||||
error: () => {},
|
||||
} as unknown as Logger
|
||||
|
||||
function makeRedis(xlen: ReturnType<typeof vi.fn>): RedisStore {
|
||||
return { xlen } as unknown as RedisStore
|
||||
}
|
||||
|
||||
function makeGate(xlen: ReturnType<typeof vi.fn>, over: Partial<{ highWater: number; lowWater: number; pollMs: number }> = {}): BackpressureGate {
|
||||
return new BackpressureGate({
|
||||
redis: makeRedis(xlen),
|
||||
stream: 's',
|
||||
highWater: over.highWater ?? 100,
|
||||
lowWater: over.lowWater ?? 50,
|
||||
pollMs: over.pollMs ?? 1,
|
||||
logger: silentLog,
|
||||
})
|
||||
}
|
||||
|
||||
describe('BackpressureGate — конструктор', () => {
|
||||
const redis = makeRedis(vi.fn())
|
||||
|
||||
it('бросает если highWater <= 0', () => {
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 0, lowWater: 0 })).toThrow(/highWater/)
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: -5, lowWater: 0 })).toThrow(/highWater/)
|
||||
})
|
||||
|
||||
it('бросает если lowWater >= highWater', () => {
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 100 })).toThrow(/lowWater/)
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 150 })).toThrow(/lowWater/)
|
||||
})
|
||||
|
||||
it('бросает если lowWater < 0', () => {
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: -1 })).toThrow(/lowWater/)
|
||||
})
|
||||
|
||||
it('принимает валидные границы', () => {
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 50 })).not.toThrow()
|
||||
expect(() => new BackpressureGate({ redis, stream: 's', highWater: 100, lowWater: 0 })).not.toThrow()
|
||||
})
|
||||
})
|
||||
|
||||
describe('BackpressureGate — waitForCapacity', () => {
|
||||
it('возвращается сразу когда backlog < highWater (одна проверка XLEN, без паузы)', async () => {
|
||||
const xlen = vi.fn().mockResolvedValue(50)
|
||||
const gate = makeGate(xlen)
|
||||
await gate.waitForCapacity()
|
||||
expect(xlen).toHaveBeenCalledTimes(1)
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
|
||||
it('возвращается сразу при backlog = highWater - 1', async () => {
|
||||
const xlen = vi.fn().mockResolvedValue(99)
|
||||
const gate = makeGate(xlen)
|
||||
await gate.waitForCapacity()
|
||||
expect(xlen).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('ставит паузу при backlog >= highWater и возобновляет при сливе до <= lowWater', async () => {
|
||||
// 150 (пауза) → опрос: 150, 70 (>50), 30 (<=50 → break)
|
||||
const xlen = vi.fn()
|
||||
.mockResolvedValueOnce(150)
|
||||
.mockResolvedValueOnce(150)
|
||||
.mockResolvedValueOnce(70)
|
||||
.mockResolvedValueOnce(30)
|
||||
const gate = makeGate(xlen)
|
||||
await gate.waitForCapacity()
|
||||
expect(xlen).toHaveBeenCalledTimes(4)
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
|
||||
it('пауза срабатывает ровно на backlog == highWater (граница >=)', async () => {
|
||||
const xlen = vi.fn().mockResolvedValueOnce(100).mockResolvedValueOnce(40)
|
||||
const gate = makeGate(xlen)
|
||||
await gate.waitForCapacity()
|
||||
expect(xlen).toHaveBeenCalledTimes(2) // 100 → пауза, 40 → выход
|
||||
})
|
||||
|
||||
it('isPaused = true во время паузы', async () => {
|
||||
const snapshots: boolean[] = []
|
||||
let gate: BackpressureGate
|
||||
const xlen = vi.fn().mockImplementation(() => {
|
||||
// снимок состояния паузы в момент опроса XLEN
|
||||
snapshots.push(gate.isPaused)
|
||||
return Promise.resolve(snapshots.length >= 3 ? 10 : 200)
|
||||
})
|
||||
gate = makeGate(xlen)
|
||||
await gate.waitForCapacity()
|
||||
// 1-й вызов — до установки паузы (false); внутри цикла — true
|
||||
expect(snapshots[0]).toBe(false)
|
||||
expect(snapshots.some(s => s === true)).toBe(true)
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
|
||||
it('shouldStop() прерывает паузу даже при высоком backlog (graceful shutdown)', async () => {
|
||||
const xlen = vi.fn().mockResolvedValue(9999) // никогда не сольётся
|
||||
const gate = makeGate(xlen)
|
||||
let stop = false
|
||||
const p = gate.waitForCapacity(() => stop)
|
||||
// через тик включаем стоп
|
||||
setTimeout(() => { stop = true }, 5)
|
||||
await p // не должно зависнуть
|
||||
expect(gate.isPaused).toBe(false)
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user