[794-5][@ant] feat: verifier-api demux-индексер permissions из native-delta parser2 — исторический суперсет authority в Postgres для доказательной point-in-time проверки ключа (C6)
This commit is contained in:
@@ -0,0 +1,47 @@
|
||||
-- E3.1 — исторический индекс permissions (append-only).
|
||||
-- Принцип «не потерять под перепарсинг»: храним ПОЛНЫЙ суперсет authority,
|
||||
-- чтобы расширение scope (owner, делегирование) НЕ требовало повторного 1–2-нед bootstrap.
|
||||
|
||||
CREATE TABLE IF NOT EXISTS permission_history (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
username TEXT NOT NULL, -- row.owner (аккаунт)
|
||||
permission_name TEXT NOT NULL, -- owner / active / кастом — БЕЗ хардкода 'active'
|
||||
parent TEXT NOT NULL DEFAULT '',
|
||||
threshold INTEGER NOT NULL DEFAULT 1,
|
||||
keys JSONB NOT NULL DEFAULT '[]'::jsonb, -- [{public_key,weight}]
|
||||
accounts JSONB NOT NULL DEFAULT '[]'::jsonb, -- [{permission:{actor,permission},weight}] — делегирование (msig, задел v2)
|
||||
waits JSONB NOT NULL DEFAULT '[]'::jsonb, -- [{wait_sec,weight}]
|
||||
block_num BIGINT NOT NULL,
|
||||
block_time TIMESTAMPTZ NOT NULL,
|
||||
is_deleted BOOLEAN NOT NULL DEFAULT FALSE, -- present=false → permission удалён в этом блоке
|
||||
event_id TEXT, -- parser2 детерминированный id (дубль-страховка)
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Идемпотентность: натуральный ключ. global_sequence для native-delta недоступен
|
||||
-- (SHiP отдаёт row-state на конец блока), поэтому (block_num, username, permission_name).
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS uq_permhist_natural
|
||||
ON permission_history (block_num, username, permission_name);
|
||||
|
||||
-- Point-in-time C6: последняя запись (username, permission_name) с block_time <= at.
|
||||
CREATE INDEX IF NOT EXISTS idx_permhist_temporal
|
||||
ON permission_history (username, permission_name, block_time DESC);
|
||||
|
||||
-- Поиск по публичному ключу (входит ли ключ в keys).
|
||||
CREATE INDEX IF NOT EXISTS idx_permhist_keys
|
||||
ON permission_history USING GIN (keys jsonb_path_ops);
|
||||
|
||||
-- Reverse-lookup делегирований (задел v2, пишем индекс сразу).
|
||||
CREATE INDEX IF NOT EXISTS idx_permhist_accounts
|
||||
ON permission_history USING GIN (accounts jsonb_path_ops);
|
||||
|
||||
-- Прогресс синхронизации индексатора + head цепи (для health-бейджа).
|
||||
CREATE TABLE IF NOT EXISTS sync_progress (
|
||||
key TEXT PRIMARY KEY,
|
||||
block_num BIGINT NOT NULL DEFAULT 0,
|
||||
block_time TIMESTAMPTZ,
|
||||
head_block_num BIGINT NOT NULL DEFAULT 0,
|
||||
head_block_time TIMESTAMPTZ,
|
||||
bootstrap_done BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
@@ -0,0 +1,55 @@
|
||||
import Fastify, { type FastifyInstance, type FastifyError } from 'fastify';
|
||||
import cors from '@fastify/cors';
|
||||
import rateLimit from '@fastify/rate-limit';
|
||||
import swagger from '@fastify/swagger';
|
||||
import { activePermissionRoutes } from './routes/active-permission.js';
|
||||
import { healthRoutes } from './routes/health.js';
|
||||
|
||||
export async function buildApp(): Promise<FastifyInstance> {
|
||||
const app = Fastify({
|
||||
logger: {
|
||||
// Pino structured-логи; усечение чувствительных полей делаем в сериализаторе.
|
||||
level: process.env.LOG_LEVEL ?? 'info',
|
||||
redact: { paths: ['req.headers.authorization'], remove: true },
|
||||
},
|
||||
trustProxy: true, // IP из X-Forwarded-For для rate-limit/логов
|
||||
});
|
||||
|
||||
await app.register(cors, {
|
||||
origin: true,
|
||||
methods: ['GET'],
|
||||
});
|
||||
|
||||
await app.register(rateLimit, {
|
||||
max: 60,
|
||||
timeWindow: '1 minute',
|
||||
// burst через allowList не нужен; 429 + Retry-After по умолчанию
|
||||
});
|
||||
|
||||
await app.register(swagger, {
|
||||
openapi: {
|
||||
openapi: '3.1.0',
|
||||
info: { title: 'Verifier API', version: '1.0.0' },
|
||||
tags: [{ name: 'active-permission' }, { name: 'health' }],
|
||||
},
|
||||
});
|
||||
|
||||
// единый формат ошибки
|
||||
app.setErrorHandler((err: FastifyError, _req, reply) => {
|
||||
const status = err.statusCode ?? 500;
|
||||
reply.code(status).send({
|
||||
error: {
|
||||
code: err.code ?? 'INTERNAL',
|
||||
message: status >= 500 ? 'Внутренняя ошибка' : err.message,
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
await app.register(activePermissionRoutes);
|
||||
await app.register(healthRoutes);
|
||||
|
||||
// отдать openapi.json
|
||||
app.get('/openapi.json', async () => app.swagger());
|
||||
|
||||
return app;
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
/** Конфиг только из env (секреты не в коде). */
|
||||
function req(name: string): string {
|
||||
const v = process.env[name];
|
||||
if (!v) throw new Error(`Не задана обязательная env-переменная ${name}`);
|
||||
return v;
|
||||
}
|
||||
|
||||
export const config = {
|
||||
databaseUrl: process.env.DATABASE_URL ?? 'postgres://verifier:verifier@localhost:5433/verifier',
|
||||
redisUrl: process.env.REDIS_URL ?? 'redis://localhost:6380',
|
||||
shipUrl: process.env.SHIP_URL ?? 'ws://127.0.0.1:8080',
|
||||
/** chain_id цепи Coopenomics — нужен ParserClient для namespacing Redis-стрима. */
|
||||
chainId: process.env.CHAIN_ID ?? '',
|
||||
startBlock: Number(process.env.START_BLOCK ?? '1'),
|
||||
port: Number(process.env.API_PORT ?? '3030'),
|
||||
host: process.env.API_HOST ?? '0.0.0.0',
|
||||
/** Порог отставания (блоков) для классификации synced/lagging. */
|
||||
laggingThreshold: Number(process.env.LAGGING_THRESHOLD ?? '50'),
|
||||
offlineThreshold: Number(process.env.OFFLINE_THRESHOLD ?? '1000'),
|
||||
get db() {
|
||||
return req('DATABASE_URL');
|
||||
},
|
||||
};
|
||||
|
||||
export const SYNC_KEY = 'permission_indexer';
|
||||
@@ -0,0 +1,14 @@
|
||||
import { Kysely, PostgresDialect } from 'kysely';
|
||||
import pg from 'pg';
|
||||
import { config } from '../config.js';
|
||||
import type { Database } from './schema.js';
|
||||
|
||||
const pool = new pg.Pool({ connectionString: config.databaseUrl, max: 10 });
|
||||
|
||||
export const db = new Kysely<Database>({
|
||||
dialect: new PostgresDialect({ pool }),
|
||||
});
|
||||
|
||||
export async function closeDb(): Promise<void> {
|
||||
await db.destroy();
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
/** Простой forward-only раннер SQL-миграций из ./migrations (идемпотентный, IF NOT EXISTS). */
|
||||
import { readdir, readFile } from 'node:fs/promises';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { dirname, join } from 'node:path';
|
||||
import pg from 'pg';
|
||||
import { config } from '../config.js';
|
||||
|
||||
const here = dirname(fileURLToPath(import.meta.url));
|
||||
// dist/db/migrate.js → ../../migrations (миграции не компилируются, лежат рядом с пакетом)
|
||||
const migrationsDir = join(here, '..', '..', 'migrations');
|
||||
|
||||
export async function migrate(): Promise<void> {
|
||||
const client = new pg.Client({ connectionString: config.databaseUrl });
|
||||
await client.connect();
|
||||
try {
|
||||
await client.query(
|
||||
'CREATE TABLE IF NOT EXISTS _migrations (name TEXT PRIMARY KEY, applied_at TIMESTAMPTZ DEFAULT now())',
|
||||
);
|
||||
const files = (await readdir(migrationsDir)).filter((f) => f.endsWith('.sql')).sort();
|
||||
for (const file of files) {
|
||||
const done = await client.query('SELECT 1 FROM _migrations WHERE name=$1', [file]);
|
||||
if (done.rowCount) continue;
|
||||
const sql = await readFile(join(migrationsDir, file), 'utf8');
|
||||
await client.query('BEGIN');
|
||||
try {
|
||||
await client.query(sql);
|
||||
await client.query('INSERT INTO _migrations(name) VALUES($1)', [file]);
|
||||
await client.query('COMMIT');
|
||||
console.log(`[migrate] applied ${file}`);
|
||||
} catch (e) {
|
||||
await client.query('ROLLBACK');
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
await client.end();
|
||||
}
|
||||
}
|
||||
|
||||
if (import.meta.url === `file://${process.argv[1]}`) {
|
||||
migrate().then(
|
||||
() => process.exit(0),
|
||||
(e) => {
|
||||
console.error(e);
|
||||
process.exit(1);
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
import type { Generated, ColumnType } from 'kysely';
|
||||
|
||||
/** JSONB authority-структуры. */
|
||||
export interface KeyWeight {
|
||||
public_key: string;
|
||||
weight: number;
|
||||
}
|
||||
export interface AccountWeight {
|
||||
permission: { actor: string; permission: string };
|
||||
weight: number;
|
||||
}
|
||||
export interface WaitWeight {
|
||||
wait_sec: number;
|
||||
weight: number;
|
||||
}
|
||||
|
||||
type Timestamp = ColumnType<Date, Date | string, Date | string>;
|
||||
|
||||
export interface PermissionHistoryTable {
|
||||
id: Generated<number>;
|
||||
username: string;
|
||||
permission_name: string;
|
||||
parent: string;
|
||||
threshold: number;
|
||||
keys: ColumnType<KeyWeight[], string, string>;
|
||||
accounts: ColumnType<AccountWeight[], string, string>;
|
||||
waits: ColumnType<WaitWeight[], string, string>;
|
||||
block_num: number;
|
||||
block_time: Timestamp;
|
||||
is_deleted: boolean;
|
||||
event_id: string | null;
|
||||
created_at: Generated<Timestamp>;
|
||||
}
|
||||
|
||||
export interface SyncProgressTable {
|
||||
key: string;
|
||||
block_num: number;
|
||||
block_time: Timestamp | null;
|
||||
head_block_num: number;
|
||||
head_block_time: Timestamp | null;
|
||||
bootstrap_done: boolean;
|
||||
updated_at: Generated<Timestamp>;
|
||||
}
|
||||
|
||||
export interface Database {
|
||||
permission_history: PermissionHistoryTable;
|
||||
sync_progress: SyncProgressTable;
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
import { buildApp } from './app.js';
|
||||
import { config } from './config.js';
|
||||
import { migrate } from './db/migrate.js';
|
||||
import { db } from './db/index.js';
|
||||
import { runIndexer } from './indexer/consumer.js';
|
||||
|
||||
async function main(): Promise<void> {
|
||||
await migrate();
|
||||
|
||||
const app = await buildApp();
|
||||
await app.listen({ port: config.port, host: config.host });
|
||||
app.log.info(`verifier-api на :${config.port}`);
|
||||
|
||||
// Индексатор запускается отдельным процессом в проде; здесь — опционально по флагу.
|
||||
if (process.env.RUN_INDEXER === '1') {
|
||||
const ac = new AbortController();
|
||||
runIndexer(db, undefined, ac.signal).catch((e) => {
|
||||
app.log.error(e, 'indexer упал');
|
||||
});
|
||||
process.on('SIGTERM', () => ac.abort());
|
||||
}
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
console.error(e);
|
||||
process.exit(1);
|
||||
});
|
||||
@@ -0,0 +1,58 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import { ParserClient, type SubscriptionFilter } from '@coopenomics/parser2';
|
||||
import type { Database } from '../db/schema.js';
|
||||
import { dispatch } from './handlers.js';
|
||||
import { indexerCurrentBlock, indexerLagBlocks } from '../metrics.js';
|
||||
import { config } from '../config.js';
|
||||
|
||||
/**
|
||||
* Потребление потока parser2 через реальный `ParserClient` (E0.1-выровнено 2026-06-03).
|
||||
*
|
||||
* Архитектура parser2: отдельный Parser-процесс читает SHiP → пишет в Redis Stream;
|
||||
* наш verifier-api — лишь consumer этого стрима. Поэтому здесь нужен только Redis + chain_id,
|
||||
* SHiP_URL относится к Parser-процессу (отдельный деплой).
|
||||
*
|
||||
* Подписка фильтрами: native-delta таблицы `permission` + fork.
|
||||
* stream() — AsyncGenerator<ParserEvent>; handler'ы sequential (FIFO per block), идемпотентны.
|
||||
*/
|
||||
export const PERMISSION_FILTERS: SubscriptionFilter[] = [
|
||||
{ kind: 'native-delta', table: 'permission' },
|
||||
{ kind: 'fork' },
|
||||
];
|
||||
|
||||
export const SUBSCRIPTION_ID = 'verifier-permission-indexer';
|
||||
|
||||
export function createPermissionClient(): ParserClient {
|
||||
if (!config.chainId) {
|
||||
throw new Error('CHAIN_ID обязателен для запуска индексатора (namespacing Redis-стрима)');
|
||||
}
|
||||
return new ParserClient({
|
||||
subscriptionId: SUBSCRIPTION_ID,
|
||||
filters: PERMISSION_FILTERS,
|
||||
startFrom: config.startBlock,
|
||||
redis: { url: config.redisUrl },
|
||||
chain: { id: config.chainId },
|
||||
noSignalHandlers: true, // SIGTERM/SIGINT обрабатываем сами в index.ts
|
||||
});
|
||||
}
|
||||
|
||||
export async function runIndexer(
|
||||
db: Kysely<Database>,
|
||||
client: ParserClient = createPermissionClient(),
|
||||
signal?: AbortSignal,
|
||||
): Promise<void> {
|
||||
signal?.addEventListener('abort', () => void client.close(), { once: true });
|
||||
try {
|
||||
for await (const ev of client.stream()) {
|
||||
if (signal?.aborted) break;
|
||||
await dispatch(db, ev);
|
||||
if (ev.kind === 'native-delta') {
|
||||
indexerCurrentBlock.set(ev.block_num);
|
||||
} else if (ev.kind === 'fork') {
|
||||
indexerLagBlocks.set(0);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
await client.close();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import type { Database } from '../db/schema.js';
|
||||
import { SYNC_KEY } from '../config.js';
|
||||
import { mapPermissionRow, type ParserEvent, type PermissionDeltaEvent } from './permission-row.js';
|
||||
|
||||
/**
|
||||
* Идемпотентный upsert строки permission по натуральному ключу
|
||||
* (block_num, username, permission_name). parser2 at-least-once → onConflict doNothing.
|
||||
*/
|
||||
export async function applyPermissionEvent(
|
||||
db: Kysely<Database>,
|
||||
ev: PermissionDeltaEvent,
|
||||
): Promise<void> {
|
||||
const row = mapPermissionRow(ev);
|
||||
await db
|
||||
.insertInto('permission_history')
|
||||
.values(row)
|
||||
.onConflict((oc) => oc.columns(['block_num', 'username', 'permission_name']).doNothing())
|
||||
.execute();
|
||||
}
|
||||
|
||||
/** Откат при реорге: удалить всё строго ПОЗЖЕ канонического блока + откатить прогресс. */
|
||||
export async function applyFork(db: Kysely<Database>, forkedFromBlock: number): Promise<void> {
|
||||
await db
|
||||
.deleteFrom('permission_history')
|
||||
.where('block_num', '>', forkedFromBlock)
|
||||
.execute();
|
||||
await db
|
||||
.updateTable('sync_progress')
|
||||
.set({ block_num: forkedFromBlock })
|
||||
.where('key', '=', SYNC_KEY)
|
||||
.where('block_num', '>', forkedFromBlock)
|
||||
.execute();
|
||||
}
|
||||
|
||||
/** Обновление прогресса синхронизации после обработки блока. */
|
||||
export async function advanceProgress(
|
||||
db: Kysely<Database>,
|
||||
blockNum: number,
|
||||
blockTime: string,
|
||||
head?: { block_num: number; block_time: string },
|
||||
): Promise<void> {
|
||||
await db
|
||||
.insertInto('sync_progress')
|
||||
.values({
|
||||
key: SYNC_KEY,
|
||||
block_num: blockNum,
|
||||
block_time: blockTime,
|
||||
head_block_num: head?.block_num ?? blockNum,
|
||||
head_block_time: head?.block_time ?? blockTime,
|
||||
bootstrap_done: false,
|
||||
})
|
||||
.onConflict((oc) =>
|
||||
oc.column('key').doUpdateSet({
|
||||
block_num: blockNum,
|
||||
block_time: blockTime,
|
||||
...(head ? { head_block_num: head.block_num, head_block_time: head.block_time } : {}),
|
||||
updated_at: new Date(),
|
||||
}),
|
||||
)
|
||||
.execute();
|
||||
}
|
||||
|
||||
/** Единая точка обработки события потока (sequential — порядок мутаций сохранён). */
|
||||
export async function dispatch(db: Kysely<Database>, ev: ParserEvent): Promise<void> {
|
||||
if (ev.kind === 'fork') {
|
||||
await applyFork(db, ev.forked_from_block);
|
||||
return;
|
||||
}
|
||||
if (ev.kind === 'native-delta' && ev.table === 'permission') {
|
||||
await applyPermissionEvent(db, ev as unknown as PermissionDeltaEvent);
|
||||
await advanceProgress(db, ev.block_num, ev.block_time);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
import { describe, it, expect } from 'vitest';
|
||||
import { mapPermissionRow, type PermissionDeltaEvent, type NativePermissionRow } from './permission-row.js';
|
||||
|
||||
function ev(row: Partial<NativePermissionRow> = {}, over: Partial<PermissionDeltaEvent> = {}): PermissionDeltaEvent {
|
||||
const data: NativePermissionRow = {
|
||||
owner: 'ant',
|
||||
name: 'active',
|
||||
parent: 'owner',
|
||||
last_updated: '2026-06-03T10:00:00.000Z',
|
||||
auth: { threshold: 1, keys: [{ key: 'EOS5aaa', weight: 1 }], accounts: [], waits: [] },
|
||||
...row,
|
||||
};
|
||||
return {
|
||||
kind: 'native-delta',
|
||||
table: 'permission',
|
||||
present: true,
|
||||
block_num: 100,
|
||||
block_time: '2026-06-03T10:00:00.000Z',
|
||||
block_id: 'blk',
|
||||
chain_id: 'chain',
|
||||
lookup_key: 'ant-active',
|
||||
event_id: 'chain:n:100:blk:permission:ant-active',
|
||||
data,
|
||||
...over,
|
||||
};
|
||||
}
|
||||
|
||||
describe('mapPermissionRow — суперсет authority (реальная форма parser2)', () => {
|
||||
it('insert active: auth.keys[].key → public_key', () => {
|
||||
const r = mapPermissionRow(ev());
|
||||
expect(r.username).toBe('ant');
|
||||
expect(r.permission_name).toBe('active');
|
||||
expect(r.is_deleted).toBe(false);
|
||||
expect(JSON.parse(r.keys)).toEqual([{ public_key: 'EOS5aaa', weight: 1 }]);
|
||||
expect(r.event_id).toBe('chain:n:100:blk:permission:ant-active');
|
||||
});
|
||||
|
||||
it('delete (present=false) → is_deleted=true', () => {
|
||||
expect(mapPermissionRow(ev({}, { present: false })).is_deleted).toBe(true);
|
||||
});
|
||||
|
||||
it('owner-строка сохраняется наравне с active', () => {
|
||||
const r = mapPermissionRow(
|
||||
ev({ name: 'owner', parent: '', auth: { threshold: 1, keys: [{ key: 'EOSowner', weight: 1 }], accounts: [], waits: [] } }),
|
||||
);
|
||||
expect(r.permission_name).toBe('owner');
|
||||
expect(JSON.parse(r.keys)).toEqual([{ public_key: 'EOSowner', weight: 1 }]);
|
||||
});
|
||||
|
||||
it('мульти-ключ + threshold', () => {
|
||||
const r = mapPermissionRow(
|
||||
ev({ auth: { threshold: 2, keys: [{ key: 'EOSa', weight: 1 }, { key: 'EOSb', weight: 1 }], accounts: [], waits: [] } }),
|
||||
);
|
||||
expect(r.threshold).toBe(2);
|
||||
expect(JSON.parse(r.keys)).toHaveLength(2);
|
||||
});
|
||||
|
||||
it('accounts-делегирование (право на аккаунт) сохраняется', () => {
|
||||
const r = mapPermissionRow(
|
||||
ev({ auth: { threshold: 1, keys: [], accounts: [{ permission: { actor: 'chairman', permission: 'active' }, weight: 1 }], waits: [] } }),
|
||||
);
|
||||
expect(JSON.parse(r.accounts)).toEqual([
|
||||
{ permission: { actor: 'chairman', permission: 'active' }, weight: 1 },
|
||||
]);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,63 @@
|
||||
import type { NativeDeltaEvent, ForkEvent, ParserEvent } from '@coopenomics/parser2';
|
||||
|
||||
/**
|
||||
* Форма строки on-chain таблицы `permission`, как её отдаёт parser2 в native-delta.
|
||||
*
|
||||
* ✅ E0.1 (2026-06-03): структура сверена с `@coopenomics/coopos-ship-reader@0.2.0`
|
||||
* (`NativePermissionRow`) через типы установленного `@coopenomics/parser2@1.1.0`.
|
||||
* Реальные ключи: `auth.keys[].key` (НЕ `public_key`), payload в `event.data`, есть `last_updated`.
|
||||
* Определяем локально (а не импортируем транзитивный ship-reader), но 1:1 с ground truth.
|
||||
*/
|
||||
export interface NativePermissionRow {
|
||||
owner: string;
|
||||
name: string;
|
||||
parent: string;
|
||||
last_updated: string;
|
||||
auth: {
|
||||
threshold: number;
|
||||
keys: Array<{ key: string; weight: number }>;
|
||||
accounts: Array<{ permission: { actor: string; permission: string }; weight: number }>;
|
||||
waits: Array<{ wait_sec: number; weight: number }>;
|
||||
};
|
||||
}
|
||||
|
||||
export type PermissionDeltaEvent = NativeDeltaEvent<NativePermissionRow>;
|
||||
export type { ParserEvent, ForkEvent };
|
||||
|
||||
/** Строка для вставки в permission_history (полный суперсет, ничего не отбрасываем). */
|
||||
export interface PermissionInsert {
|
||||
username: string;
|
||||
permission_name: string;
|
||||
parent: string;
|
||||
threshold: number;
|
||||
keys: string;
|
||||
accounts: string;
|
||||
waits: string;
|
||||
block_num: number;
|
||||
block_time: string;
|
||||
is_deleted: boolean;
|
||||
event_id: string | null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Маппинг native-delta события `permission` → строка суперсета.
|
||||
* `auth.keys[].key` нормализуем в `{public_key, weight}` (имя поля для downstream C6-запросов).
|
||||
*/
|
||||
export function mapPermissionRow(ev: PermissionDeltaEvent): PermissionInsert {
|
||||
const row = ev.data;
|
||||
return {
|
||||
username: row.owner,
|
||||
permission_name: row.name,
|
||||
parent: row.parent ?? '',
|
||||
threshold: row.auth?.threshold ?? 1,
|
||||
keys: JSON.stringify(
|
||||
(row.auth?.keys ?? []).map((k) => ({ public_key: k.key, weight: k.weight })),
|
||||
),
|
||||
accounts: JSON.stringify(row.auth?.accounts ?? []),
|
||||
waits: JSON.stringify(row.auth?.waits ?? []),
|
||||
block_num: ev.block_num,
|
||||
block_time: ev.block_time,
|
||||
is_deleted: ev.present === false,
|
||||
event_id: ev.event_id ?? null,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
import { Registry, Gauge, collectDefaultMetrics } from 'prom-client';
|
||||
|
||||
export const registry = new Registry();
|
||||
collectDefaultMetrics({ register: registry });
|
||||
|
||||
export const indexerCurrentBlock = new Gauge({
|
||||
name: 'indexer_current_block',
|
||||
help: 'Последний обработанный индексатором блок',
|
||||
registers: [registry],
|
||||
});
|
||||
|
||||
export const indexerLagBlocks = new Gauge({
|
||||
name: 'indexer_lag_blocks',
|
||||
help: 'Отставание индексатора от head цепи (в блоках)',
|
||||
registers: [registry],
|
||||
});
|
||||
|
||||
export const indexerBootstrapDone = new Gauge({
|
||||
name: 'indexer_bootstrap_done',
|
||||
help: '1 — первичная индексация (bootstrap) завершена',
|
||||
registers: [registry],
|
||||
});
|
||||
@@ -0,0 +1,120 @@
|
||||
import type { Kysely } from 'kysely';
|
||||
import type { Database, KeyWeight } from './db/schema.js';
|
||||
import { SYNC_KEY } from './config.js';
|
||||
import { classify, type HealthView } from './services/health-status.js';
|
||||
|
||||
export interface CheckResultRow {
|
||||
active: boolean;
|
||||
last_known_block: number | null;
|
||||
last_known_block_time: string | null;
|
||||
/** at позже последнего проиндексированного блока — «белая зона». */
|
||||
beyond_index: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* C6 point-in-time: был ли public_key в active-permission аккаунта на момент `at`.
|
||||
* Берём последнюю запись (username, 'active') c block_time <= at (по idx_permhist_temporal),
|
||||
* проверяем present (NOT is_deleted) и вхождение ключа в keys.
|
||||
*/
|
||||
export async function checkActiveKey(
|
||||
db: Kysely<Database>,
|
||||
username: string,
|
||||
publicKey: string,
|
||||
at: string,
|
||||
): Promise<CheckResultRow | null> {
|
||||
const last = await db
|
||||
.selectFrom('permission_history')
|
||||
.selectAll()
|
||||
.where('username', '=', username)
|
||||
.where('permission_name', '=', 'active')
|
||||
.where('block_time', '<=', new Date(at))
|
||||
.orderBy('block_time', 'desc')
|
||||
.orderBy('block_num', 'desc')
|
||||
.limit(1)
|
||||
.executeTakeFirst();
|
||||
|
||||
// существует ли аккаунт в индексе вообще?
|
||||
if (!last) {
|
||||
const any = await db
|
||||
.selectFrom('permission_history')
|
||||
.select('id')
|
||||
.where('username', '=', username)
|
||||
.limit(1)
|
||||
.executeTakeFirst();
|
||||
if (!any) return null; // 404 — аккаунта нет в индексе
|
||||
}
|
||||
|
||||
const progress = await getHealth(db);
|
||||
const beyond_index = progress.current_block_time !== null && at > progress.current_block_time;
|
||||
|
||||
if (!last) {
|
||||
return { active: false, last_known_block: null, last_known_block_time: null, beyond_index };
|
||||
}
|
||||
|
||||
const keys = last.keys as unknown as KeyWeight[];
|
||||
const active = !last.is_deleted && keys.some((k) => k.public_key === publicKey);
|
||||
return {
|
||||
active,
|
||||
last_known_block: last.block_num,
|
||||
last_known_block_time: toIso(last.block_time),
|
||||
beyond_index,
|
||||
};
|
||||
}
|
||||
|
||||
/** История изменений active-permission аккаунта (по возрастанию времени). */
|
||||
export async function getActiveHistory(db: Kysely<Database>, username: string) {
|
||||
const rows = await db
|
||||
.selectFrom('permission_history')
|
||||
.select(['block_num', 'block_time', 'keys', 'threshold', 'is_deleted', 'event_id'])
|
||||
.where('username', '=', username)
|
||||
.where('permission_name', '=', 'active')
|
||||
.orderBy('block_time', 'asc')
|
||||
.orderBy('block_num', 'asc')
|
||||
.execute();
|
||||
return rows.map((r) => ({
|
||||
block_num: r.block_num,
|
||||
block_time: toIso(r.block_time),
|
||||
keys: r.keys as unknown as KeyWeight[],
|
||||
threshold: r.threshold,
|
||||
present: !r.is_deleted,
|
||||
event_id: r.event_id,
|
||||
}));
|
||||
}
|
||||
|
||||
export async function getHealth(db: Kysely<Database>): Promise<HealthView> {
|
||||
const p = await db
|
||||
.selectFrom('sync_progress')
|
||||
.selectAll()
|
||||
.where('key', '=', SYNC_KEY)
|
||||
.executeTakeFirst();
|
||||
|
||||
if (!p) {
|
||||
return {
|
||||
head_block: 0,
|
||||
head_block_time: null,
|
||||
current_block: 0,
|
||||
current_block_time: null,
|
||||
behind_by: 0,
|
||||
status: 'offline',
|
||||
bootstrap_done: false,
|
||||
source: 'verifier-indexer',
|
||||
};
|
||||
}
|
||||
|
||||
const behind_by = Math.max(0, Number(p.head_block_num) - Number(p.block_num));
|
||||
return {
|
||||
head_block: Number(p.head_block_num),
|
||||
head_block_time: toIso(p.head_block_time),
|
||||
current_block: Number(p.block_num),
|
||||
current_block_time: toIso(p.block_time),
|
||||
behind_by,
|
||||
status: classify(behind_by, true),
|
||||
bootstrap_done: p.bootstrap_done,
|
||||
source: 'verifier-indexer',
|
||||
};
|
||||
}
|
||||
|
||||
function toIso(v: Date | string | null): string | null {
|
||||
if (v === null) return null;
|
||||
return v instanceof Date ? v.toISOString() : new Date(v).toISOString();
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
import type { FastifyInstance } from 'fastify';
|
||||
import { Type } from '@sinclair/typebox';
|
||||
import { db } from '../db/index.js';
|
||||
import { checkActiveKey, getActiveHistory } from '../repo.js';
|
||||
|
||||
const CheckQuery = Type.Object({
|
||||
username: Type.String({ minLength: 1, maxLength: 13 }),
|
||||
public_key: Type.String({ minLength: 8 }),
|
||||
at: Type.String({ format: 'date-time' }),
|
||||
});
|
||||
|
||||
const HistoryQuery = Type.Object({
|
||||
username: Type.String({ minLength: 1, maxLength: 13 }),
|
||||
});
|
||||
|
||||
export async function activePermissionRoutes(app: FastifyInstance): Promise<void> {
|
||||
// C6 point-in-time
|
||||
app.get(
|
||||
'/v1/active-permission/check',
|
||||
{
|
||||
schema: {
|
||||
querystring: CheckQuery,
|
||||
description: 'Был ли public_key активен (active-permission) у username на момент at',
|
||||
tags: ['active-permission'],
|
||||
},
|
||||
},
|
||||
async (req, reply) => {
|
||||
const { username, public_key, at } = req.query as {
|
||||
username: string;
|
||||
public_key: string;
|
||||
at: string;
|
||||
};
|
||||
const res = await checkActiveKey(db, username, public_key, at);
|
||||
if (res === null) {
|
||||
return reply.code(404).send({
|
||||
error: { code: 'ACCOUNT_NOT_INDEXED', message: `Аккаунт ${username} не найден в индексе` },
|
||||
});
|
||||
}
|
||||
if (res.beyond_index) {
|
||||
return reply.code(503).send({
|
||||
error: {
|
||||
code: 'BEYOND_INDEX',
|
||||
message: 'Момент at позже последнего проиндексированного блока (белая зона)',
|
||||
details: { synced_to: res.last_known_block_time },
|
||||
},
|
||||
});
|
||||
}
|
||||
return {
|
||||
active: res.active,
|
||||
username,
|
||||
public_key,
|
||||
at,
|
||||
scope: 'active-permission',
|
||||
last_known_block: res.last_known_block,
|
||||
last_known_block_time: res.last_known_block_time,
|
||||
source: 'verifier-indexer',
|
||||
};
|
||||
},
|
||||
);
|
||||
|
||||
// история active-permission
|
||||
app.get(
|
||||
'/v1/active-permission/history',
|
||||
{
|
||||
schema: {
|
||||
querystring: HistoryQuery,
|
||||
description: 'История изменений active-permission аккаунта',
|
||||
tags: ['active-permission'],
|
||||
},
|
||||
},
|
||||
async (req) => {
|
||||
const { username } = req.query as { username: string };
|
||||
return { username, scope: 'active-permission', events: await getActiveHistory(db, username) };
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
import type { FastifyInstance } from 'fastify';
|
||||
import { db } from '../db/index.js';
|
||||
import { getHealth } from '../repo.js';
|
||||
import { registry } from '../metrics.js';
|
||||
|
||||
export async function healthRoutes(app: FastifyInstance): Promise<void> {
|
||||
app.get(
|
||||
'/v1/health',
|
||||
{ schema: { description: 'Состояние demux-индексатора', tags: ['health'] } },
|
||||
async () => getHealth(db),
|
||||
);
|
||||
|
||||
app.get('/metrics', async (_req, reply) => {
|
||||
reply.header('Content-Type', registry.contentType);
|
||||
return registry.metrics();
|
||||
});
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
import { config } from '../config.js';
|
||||
|
||||
export type SyncStatus = 'synced' | 'lagging' | 'offline';
|
||||
|
||||
export interface HealthView {
|
||||
head_block: number;
|
||||
head_block_time: string | null;
|
||||
current_block: number;
|
||||
current_block_time: string | null;
|
||||
behind_by: number;
|
||||
status: SyncStatus;
|
||||
bootstrap_done: boolean;
|
||||
source: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Классификация состояния индексатора по отставанию (правила §6 архитектуры):
|
||||
* synced — behind_by <= laggingThreshold
|
||||
* lagging — laggingThreshold < behind_by <= offlineThreshold
|
||||
* offline — behind_by > offlineThreshold (или нет данных)
|
||||
*/
|
||||
export function classify(behindBy: number, hasData: boolean): SyncStatus {
|
||||
if (!hasData) return 'offline';
|
||||
if (behindBy <= config.laggingThreshold) return 'synced';
|
||||
if (behindBy <= config.offlineThreshold) return 'lagging';
|
||||
return 'offline';
|
||||
}
|
||||
Reference in New Issue
Block a user