Merge pull request '[E10-3b] orchestrator: dynamic supergraph manager — recompose без рестарта' (#75) from feat/E10-3b-dynamic-supergraph into feat/E10-5-on-chain-watcher

This commit was merged in pull request #75.
This commit is contained in:
2026-06-04 11:23:04 +00:00
2 changed files with 305 additions and 0 deletions
@@ -0,0 +1,183 @@
/**
* @fileoverview Юнит-тесты dynamic supergraph manager'а (Story 10.3b).
*
* Покрытие:
* 1. initial compose: возвращает SDL, composer вызывается ровно один раз;
* 2. tick без изменений в registry → composer НЕ вызывается, update НЕ вызывается;
* 3. tick с новой subgraph-записью в registry → composer вызывается, update вызывается с новым SDL;
* 4. tick с удалённой subgraph-записью → recompose;
* 5. tick с изменённым URL у subgraph'а (рестарт сервиса на другом порту) → recompose;
* 6. ошибка в composer'е на tick'е НЕ ломает следующий tick (continuity);
* 7. cleanup() останавливает таймер — composer перестаёт вызываться.
*
* Используем jest fake timers для управления интервалом.
*/
import {
createDynamicSupergraphManager,
SupergraphComposerPort,
SupergraphRegistryReader,
} from './supergraph-manager';
import type { SubgraphDescriptor } from './subgraph-registry.service';
class FakeRegistry implements SupergraphRegistryReader {
current: SubgraphDescriptor[] = [];
calls = 0;
async listForCompose(): Promise<SubgraphDescriptor[]> {
this.calls += 1;
return this.current.slice();
}
}
class FakeComposer implements SupergraphComposerPort {
calls: Array<ReadonlyArray<SubgraphDescriptor>> = [];
results: string[] = [];
error?: Error;
async compose(subgraphs: ReadonlyArray<SubgraphDescriptor>): Promise<string> {
this.calls.push(subgraphs.slice());
if (this.error) throw this.error;
const sdl = `# supergraph v${this.calls.length}\n${subgraphs.map((s) => s.name).join(',')}`;
this.results.push(sdl);
return sdl;
}
}
const flushMicrotasks = async (): Promise<void> => {
await Promise.resolve();
await Promise.resolve();
await Promise.resolve();
};
const POLL_MS = 1000;
describe('createDynamicSupergraphManager', () => {
beforeEach(() => {
jest.useFakeTimers();
});
afterEach(() => {
jest.useRealTimers();
});
it('initial compose: возвращает SDL, composer вызывается ровно один раз', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
expect(composer.calls.length).toBe(1);
expect(lifecycle.initialSdl).toBe('# supergraph v1\ncore');
expect(updates).toEqual([]);
await lifecycle.cleanup();
});
it('tick без изменений → composer НЕ вызывается, update НЕ вызывается', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
expect(composer.calls.length).toBe(1);
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(composer.calls.length).toBe(1); // не вызван повторно
expect(updates).toEqual([]);
await lifecycle.cleanup();
});
it('tick с новой subgraph-записью → composer вызывается, update получает новый SDL', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
registry.current.push({ name: 'chatcoop', url: 'http://chatcoop:3000/graphql' });
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(composer.calls.length).toBe(2);
expect(updates).toEqual(['# supergraph v2\ncore,chatcoop']);
await lifecycle.cleanup();
});
it('tick с удалённой subgraph-записью → recompose', async () => {
const registry = new FakeRegistry();
registry.current = [
{ name: 'core', url: 'http://core:3000/graphql' },
{ name: 'chatcoop', url: 'http://chatcoop:3000/graphql' },
];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(composer.calls.length).toBe(2);
expect(updates).toEqual(['# supergraph v2\ncore']);
await lifecycle.cleanup();
});
it('tick с изменённым URL у subgraph\'а → recompose', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'chatcoop', url: 'http://chatcoop:3000/graphql' }];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
registry.current = [{ name: 'chatcoop', url: 'http://chatcoop-v2:3000/graphql' }];
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(composer.calls.length).toBe(2);
expect(updates.length).toBe(1);
await lifecycle.cleanup();
});
it('ошибка в composer\'е НЕ ломает следующий tick (continuity)', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
const composer = new FakeComposer();
const updates: string[] = [];
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
(sdl) => updates.push(sdl),
);
// tick 1 — добавили subgraph, composer бросает
registry.current.push({ name: 'chatcoop', url: 'http://chatcoop:3000/graphql' });
composer.error = new Error('introspection failed');
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(updates).toEqual([]);
// tick 2 — composer оправился
composer.error = undefined;
jest.advanceTimersByTime(POLL_MS);
await flushMicrotasks();
expect(updates.length).toBe(1);
await lifecycle.cleanup();
});
it('cleanup() останавливает таймер — composer перестаёт вызываться', async () => {
const registry = new FakeRegistry();
registry.current = [{ name: 'core', url: 'http://core:3000/graphql' }];
const composer = new FakeComposer();
const lifecycle = await createDynamicSupergraphManager(
{ registry, composer, pollIntervalMs: POLL_MS },
() => undefined,
);
await lifecycle.cleanup();
registry.current.push({ name: 'chatcoop', url: 'http://chatcoop:3000/graphql' });
jest.advanceTimersByTime(POLL_MS * 5);
await flushMicrotasks();
expect(composer.calls.length).toBe(1); // только initial
});
});
@@ -0,0 +1,122 @@
/**
* @fileoverview Custom SupergraphManager — Story 10.3b.
*
* Apollo Gateway по умолчанию через `IntrospectAndCompose` принимает
* СТАТИЧЕСКИЙ список subgraph'ов на bootstrap'е. Polling-режим
* подхватывает изменения СХЕМЫ у уже известных subgraph'ов, но НЕ
* замечает появление НОВЫХ записей в registry — для них требуется
* рестарт контейнера.
*
* Story 10.4 install pipeline пишет новые subgraph'ы в registry без
* рестарта; чтобы gateway их видел, нужен SupergraphManager, который
* на каждый tick re-читает registry, и если список (по `(name,url)`)
* изменился — пересобирает supergraph и вызывает `update(newSdl)`.
*
* Apollo Gateway это поддерживает: `gateway: { supergraphSdl: fn }`,
* где fn возвращает `{ supergraphSdl, cleanup }` и получает callback
* `update`. См. https://www.apollographql.com/docs/federation/v2/api/apollo-gateway/#supergraphsdl
*
* Этот файл — переиспользуемый менеджер, отделённый от Apollo Gateway
* API (его hook-fn в `gateway.module.ts` собирает менеджер и
* передаёт в gateway). Сам менеджер не знает про gateway — он
* работает с двумя портами: regsitry (читать список) и
* supergraphComposer (собрать SDL из subgraph descriptors).
*/
import { Logger } from '@nestjs/common';
import type { SubgraphDescriptor } from './subgraph-registry.service';
/**
* Порт композитора supergraph'а. Реальный impl делает introspection
* subgraph URL'ов и compose через `@apollo/composition`. Тестовый
* возвращает заранее заданный SDL.
*/
export interface SupergraphComposerPort {
compose(subgraphs: ReadonlyArray<SubgraphDescriptor>): Promise<string>;
}
/**
* Минимальный API регистра, который нужен менеджеру.
* Преднамеренно ужe, чем весь `SubgraphRegistryService` — менеджер
* вообще не должен знать про писать/изменять записи.
*/
export interface SupergraphRegistryReader {
listForCompose(): Promise<SubgraphDescriptor[]>;
}
export interface SupergraphManagerOptions {
composer: SupergraphComposerPort;
registry: SupergraphRegistryReader;
/** Интервал опроса registry. Должен совпадать с polling'ом IntrospectAndCompose. */
pollIntervalMs: number;
}
export interface SupergraphManagerLifecycle {
/** Начальный SDL для отдачи gateway'ю на bootstrap'е. */
initialSdl: string;
/** Освободить таймер — gateway вызовет при shutdown'е. */
cleanup(): Promise<void>;
}
/**
* Создаёт менеджер supergraph'а с динамическим refresh'ем по registry.
*
* 1. Делает первый compose(registry.listForCompose()) → возвращает
* `initialSdl` и сохраняет «текущее состояние» (отпечаток списка).
* 2. Запускает setInterval, который на каждый tick читает registry;
* если отпечаток списка отличается от текущего — recompose и
* вызов `update(newSdl)`.
* 3. `cleanup()` останавливает таймер.
*
* Идемпотентность: если registry вернул тот же список (по name+url) —
* compose НЕ вызывается, экономим CPU и сеть (introspection).
*/
export async function createDynamicSupergraphManager(
opts: SupergraphManagerOptions,
update: (sdl: string) => void,
): Promise<SupergraphManagerLifecycle> {
const logger = new Logger('DynamicSupergraphManager');
let lastFingerprint = '';
const fingerprintOf = (subgraphs: ReadonlyArray<SubgraphDescriptor>): string =>
subgraphs
.slice()
.sort((a, b) => a.name.localeCompare(b.name))
.map((s) => `${s.name}:${s.url}`)
.join('|');
const composeOnce = async (): Promise<string> => {
const subgraphs = await opts.registry.listForCompose();
const fp = fingerprintOf(subgraphs);
lastFingerprint = fp;
return opts.composer.compose(subgraphs);
};
const initialSdl = await composeOnce();
const timer = setInterval(() => {
void (async () => {
try {
const subgraphs = await opts.registry.listForCompose();
const fp = fingerprintOf(subgraphs);
if (fp === lastFingerprint) return;
logger.log(`supergraph registry changed (was ${lastFingerprint.split('|').length} subgraphs, now ${fp.split('|').length}) — recompose`);
// ВАЖНО: fingerprint обновляем только ПОСЛЕ успешного compose.
// Иначе при exception в composer'е следующий tick посчитает,
// что состояние уже актуально (fp === lastFingerprint) и пропустит
// retry — supergraph навсегда останется в устаревшем состоянии.
const sdl = await opts.composer.compose(subgraphs);
lastFingerprint = fp;
update(sdl);
} catch (e) {
logger.error(`recompose tick failed: ${e instanceof Error ? e.message : String(e)}`);
}
})();
}, opts.pollIntervalMs);
return {
initialSdl,
cleanup: async () => {
clearInterval(timer);
},
};
}