Шардировать пакетную миграцию схем тенантов

Обход реестра идёт страницами (ITenantRepository.ListPageAsync, DefaultPageSize=200) с параллелизмом внутри страницы и изоляцией сбоев — масштаб на 1000+ схем без загрузки всего реестра. Закрывает BL-SCALE-1000.
This commit is contained in:
Rustam Khalimov
2026-09-11 17:43:27 +03:00
parent 38561528ab
commit 1d7c9980ba
12 changed files with 273 additions and 44 deletions
+1 -1
View File
@@ -36,7 +36,7 @@
|---|---|---|---|---|
| BL-BILLING | Биллинг и тарифные планы, провайдер платежей | ТЗ §12 | P3 | DEFERRED |
| BL-SIGNUP | Саморегистрация тенантов (сейчас инвайты/оператор) | ТЗ §12 | P3 | DEFERRED |
| BL-SCALE-1000 | Механизм миграций/провижининга на 1000+ схем (пакетная миграция уже есть; шардирование обхода) | roadmap этап 0, этап 12 | P2 | BACKLOG |
| BL-SCALE-1000 | Механизм миграций/провижининга на 1000+ схем. **Сделано (2026-09-11):** пакетная миграция шардирована — `ITenantRepository.ListPageAsync` + обход страницами в `TenantSchemaMigrationService` (параллелизм внутри страницы, `DefaultPageSize=200`, границы 1..32 / 1..5000), сбои изолированы. Осталось при росте: вынести параллелизм/размер в конфиг и кэш прогресса (при необходимости) | roadmap этап 0, этап 12 | P2 | DONE |
## 4. Безопасность и наблюдаемость (доработки)
+1 -1
View File
@@ -455,7 +455,7 @@ links/files/history/tzText/reminder` — модули; `isNew/prevCol/isVacancy/
| `GET /operator/analytics/tokens?groupBy=&tenantId=&from=&to=` | агрегаты расхода токенов (`groupBy=day\|tenant\|provider\|model`) | `{groupBy,from,to,items:[{key,…}],total}`; 400 (неизвестная группировка) |
| `GET /operator/analytics/activity?eventType=&actorType=&actorId=&tenantId=&from=&to=&limit=&offset=` | лента действий (аудит) с фильтрами и пагинацией | `{items,total,limit,offset}` |
| `GET /operator/health` | health core/БД + сервисы ml/ai/telegram (UseLocal → `mode:local`); этап 12: глубины очередей и активные сессии | `{ok, core:{db}, services:[…], queues:{pipeline,mlOutbox}, sessions:{active}}` (всегда 200) |
| `POST /operator/maintenance/tenants/migrate` | Пакетная миграция схем всех тенантов (идемпотентно, ограниченный параллелизм; этап 12, пакет C) | `{ok,total,migrated,failed,failedSchemas,durationMs}` (`ok=false`, если хотя бы одна схема не мигрирована); 401 без операторской сессии |
| `POST /operator/maintenance/tenants/migrate` | Пакетная миграция схем всех тенантов (идемпотентно, шардированный обход страницами + ограниченный параллелизм; этап 12, пакет C / BL-SCALE-1000) | `{ok,total,migrated,failed,failedSchemas,durationMs}` (`ok=false`, если хотя бы одна схема не мигрирована); 401 без операторской сессии |
| `GET /operator/analytics/suspicious?from=&to=` | подозрительная активность по аудиту (всплеск неудачных входов по IP/логину, входы актора с множества IP, серии по тенанту; этап 12) | `{scanned,truncated,items:[…]}` |
| `GET /operator/settings/telegram-keys` | глобальные ключи Telegram (задаёт оператор; тенант их не видит) | `{apiId, apiHash (маска), keysSet}` |
| `PUT /operator/settings/telegram-keys` `{apiId?, apiHash?}` | задать/обновить ключи (частично: можно одно поле, второе сохраняется); `api_id` 59 цифр, `api_hash` непустой; hash шифруется | маска-форма; 400 `{detail}`; 401 |
+3 -2
View File
@@ -12,8 +12,9 @@
> `SourceIngressGrpcService`), `PushMessage` из telegram.proto удалён. Сухой прогон текста по конвейеру
> (стоп-правила → ML → ИИ) без записи: `POST /api/admin/check-message` + UI настроек. Remote-просмотр
> исходника: `TelegramService.ReadSource` + `TelegramSourceContentProvider` + UI «Обновить из источника».
> Метрика алертинга `deal.ai.budget.used.ratio{tenant}`. Ядро: build 5 sln 0/0,
> `Deal.Tests.Unit` **1310/1310 PASS**, telegram **130/130**, фронт `build` + `lint:i18n` зелёные.
> Метрика алертинга `deal.ai.budget.used.ratio{tenant}`. Hardening контейнеров (non-root/read-only/limits),
> шардированная пакетная миграция схем. Ядро: build 5 sln 0/0, `Deal.Tests.Unit` **1315/1315 PASS**,
> telegram **130/130**, фронт `build` + `lint:i18n` зелёные.
> Детали — `docs/superpowers/specs/2026-09-11-source-contract-design.md`.
> Осталось (в backlog): `GET /api/cards/{id}/source` + `ISourceContentProvider`, выгрузка вложений
> telegram-адаптером в Storage, `TelegramSourceContentProvider`, перенос оставшейся Telegram-специфики
@@ -636,7 +636,9 @@ DEAL_MTLS_ENABLED=0|1 DEAL_MTLS_CERT_PASSWORD=... DEAL_DEFAULT_AI_BU
- Карта `/api` — `docs/api/api-map.md` + контракты `docs/architecture/2026-09-10-unified-api-contract.md`
и `docs/architecture/2026-09-10-operator-analytics-contract.md` (актуальны на этап 12).
- Пакетная миграция схем тенантов (сотни/тысячи) — реализована на этапе 12 (`POST
/api/operator/maintenance/tenants/migrate`, §13.10/§16; см. также §4/§7).
/api/operator/maintenance/tenants/migrate`, §13.10/§16; см. также §4/§7); с BL-SCALE-1000 (2026-09-11)
обход шардирован страницами (`ITenantRepository.ListPageAsync`, `DefaultPageSize=200`) с параллелизмом
внутри страницы и изоляцией сбоев.
- Kafka — отложена.
---
@@ -38,6 +38,24 @@ public sealed class TenantRepository(DealDbContext dbContext) : ITenantRepositor
var entities = await dbContext.Tenants
.AsNoTracking()
.OrderBy(t => t.CreatedAt)
.ThenBy(t => t.Id)
.Select(t => ToTenantRecordDto(t))
.ToListAsync(ct);
return entities;
}
/// <inheritdoc />
public async Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct)
{
var entities = await dbContext.Tenants
.AsNoTracking()
.OrderBy(t => t.CreatedAt)
.ThenBy(t => t.Id)
.Skip(offset)
.Take(limit)
.Select(t => ToTenantRecordDto(t))
.ToListAsync(ct);
return entities;
@@ -8,7 +8,8 @@ using Microsoft.Extensions.Logging;
namespace Deal.Infrastructure.Tenancy;
/// <summary>
/// Пакетная (maintenance) миграция схем ВСЕХ существующих тенантов реестра с ограниченным параллелизмом и логированием прогресса — для провижининга SaaS на сотни/тысячи схем.
/// Пакетная (maintenance) миграция схем ВСЕХ существующих тенантов реестра: шардированный обход
/// страницами с ограниченным параллелизмом и логированием прогресса — для провижининга на сотни/тысячи схем.
/// </summary>
public sealed class TenantSchemaMigrationService(
ITenantRepository tenantRepository,
@@ -20,42 +21,72 @@ public sealed class TenantSchemaMigrationService(
/// </summary>
public const int DefaultMaxParallelism = 4;
/// <summary>
/// Размер страницы обхода реестра по умолчанию (шард).
/// </summary>
public const int DefaultPageSize = 200;
// Нижняя граница параллелизма (ноль/отрицательное значение недопустимо).
private const int MinMaxParallelism = 1;
// Верхняя граница параллелизма: защита Postgres от лавины одновременных DDL-подключений.
private const int MaxMaxParallelism = 32;
private const int MinPageSize = 1;
private const int MaxPageSize = 5000;
/// <summary>
/// Мигрирует схемы всех тенантов реестра с параллелизмом по умолчанию.
/// Мигрирует схемы всех тенантов реестра с параметрами по умолчанию.
/// </summary>
/// <returns>Итоговая сводка пакетной миграции.</returns>
public Task<TenantMigrationSummary> MigrateAllAsync(CancellationToken ct)
=> MigrateAllAsync(DefaultMaxParallelism, ct);
=> MigrateAllAsync(DefaultMaxParallelism, DefaultPageSize, ct);
/// <summary>
/// Мигрирует схемы всех тенантов реестра с заданным параллелизмом
/// Мигрирует схемы всех тенантов реестра с заданным параллелизмом.
/// </summary>
/// <param name="maxParallelism">Желаемое число схем, мигрируемых одновременно.</param>
/// <returns>Итоговая сводка пакетной миграции.</returns>
public async Task<TenantMigrationSummary> MigrateAllAsync(int maxParallelism, CancellationToken ct)
public Task<TenantMigrationSummary> MigrateAllAsync(int maxParallelism, CancellationToken ct)
=> MigrateAllAsync(maxParallelism, DefaultPageSize, ct);
/// <summary>
/// Мигрирует схемы всех тенантов реестра шардами по страницам.
/// </summary>
/// <param name="maxParallelism">Желаемое число схем, мигрируемых одновременно.</param>
/// <param name="pageSize">Размер страницы (шарда) обхода реестра.</param>
/// <returns>Итоговая сводка пакетной миграции.</returns>
public async Task<TenantMigrationSummary> MigrateAllAsync(
int maxParallelism,
int pageSize,
CancellationToken ct)
{
int parallelism = Math.Clamp(maxParallelism, MinMaxParallelism, MaxMaxParallelism);
IReadOnlyList<TenantRecordDto> tenants = await tenantRepository.ListAsync(ct).ConfigureAwait(false);
int shardSize = Math.Clamp(pageSize, MinPageSize, MaxPageSize);
var stopwatch = Stopwatch.StartNew();
var failedSchemas = new ConcurrentBag<string>();
int total = 0;
int migrated = 0;
int shardIndex = 0;
if (tenants.Count == 0)
while (true)
{
stopwatch.Stop();
logger.LogInformation("Пакетная миграция схем: в реестре нет тенантов — мигрировать нечего");
return new TenantMigrationSummary(0, 0, 0, stopwatch.ElapsedMilliseconds, []);
ct.ThrowIfCancellationRequested();
IReadOnlyList<TenantRecordDto> shard = await tenantRepository
.ListPageAsync(total, shardSize, ct)
.ConfigureAwait(false);
if (shard.Count == 0)
{
break;
}
var failedSchemas = new ConcurrentBag<string>();
int migrated = 0;
shardIndex++;
total += shard.Count;
int processedSoFar = total;
await Parallel.ForEachAsync(
tenants,
shard,
new ParallelOptions { MaxDegreeOfParallelism = parallelism, CancellationToken = ct },
async (tenant, token) =>
{
@@ -65,9 +96,10 @@ public sealed class TenantSchemaMigrationService(
await tenantProvisioner.ProvisionAsync(tenantId, token).ConfigureAwait(false);
int done = Interlocked.Increment(ref migrated);
logger.LogInformation(
"Пакетная миграция схем: {Done}/{Total} — {Schema} готова",
"Пакетная миграция схем: шард {Shard}, {Done}/{Total} — {Schema} готова",
shardIndex,
done,
tenants.Count,
processedSoFar,
tenantId.SchemaName);
}
catch (Exception exception) when (exception is not OperationCanceledException)
@@ -75,14 +107,28 @@ public sealed class TenantSchemaMigrationService(
failedSchemas.Add(tenantId.SchemaName);
logger.LogError(
exception,
"Пакетная миграция схем: {Schema} не мигрирована",
"Пакетная миграция схем: шард {Shard}, {Schema} не мигрирована",
shardIndex,
tenantId.SchemaName);
}
}).ConfigureAwait(false);
if (shard.Count < shardSize)
{
break;
}
}
stopwatch.Stop();
if (total == 0)
{
logger.LogInformation("Пакетная миграция схем: в реестре нет тенантов — мигрировать нечего");
return new TenantMigrationSummary(0, 0, 0, stopwatch.ElapsedMilliseconds, []);
}
var summary = new TenantMigrationSummary(
tenants.Count,
total,
migrated,
failedSchemas.Count,
stopwatch.ElapsedMilliseconds,
@@ -26,6 +26,17 @@ public interface ITenantRepository
/// <returns>Список тенантов.</returns>
public Task<IReadOnlyList<TenantRecordDto>> ListAsync(CancellationToken ct);
/// <summary>
/// Страница реестра тенантов (шардированный обход для 1000+ схем).
/// </summary>
/// <param name="offset">Сдвиг от начала (устойчивый порядок — CreatedAt, затем Id).</param>
/// <param name="limit">Размер страницы (≥1; валидирует потребитель).</param>
/// <returns>Записи страницы; пусто — страниц больше нет.</returns>
public Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct);
/// <summary>
/// Устанавливает статус тенанта.
/// </summary>
@@ -71,6 +71,98 @@ public sealed class TenantSchemaMigrationServiceTests
Assert.Equal(new[] { $"tenant_{TenantB:N}" }, provisioner.ProvisionedSchemaNames);
}
[Fact]
public async Task MigrateAllAsync_MultiplePages_ReadsEveryTenant()
{
var provisioner = new FakeTenantProvisioner();
var repository = new FakeTenantRepository(Records(5));
var service = NewService(repository, provisioner);
TenantMigrationSummary summary = await service.MigrateAllAsync(4, 2, CancellationToken.None);
Assert.True(summary.Ok);
Assert.Equal(5, summary.Total);
Assert.Equal(5, summary.Migrated);
Assert.Equal(0, summary.Failed);
Assert.Equal(5, provisioner.ProvisionedSchemaNames.Count);
Assert.Equal(5, provisioner.ProvisionedSchemaNames.Distinct().Count());
// Страницы запрошены последовательно с устойчивым offset: 0, 2, 4 (страница на 4 неполная — обход остановлен).
Assert.Equal(
new (int Offset, int Limit)[] { (0, 2), (2, 2), (4, 2) },
repository.PageRequests);
}
[Fact]
public async Task MigrateAllAsync_PageSizeOne_ReadsEveryTenant()
{
var provisioner = new FakeTenantProvisioner();
var repository = new FakeTenantRepository(Records(4));
var service = NewService(repository, provisioner);
TenantMigrationSummary summary = await service.MigrateAllAsync(2, 1, CancellationToken.None);
Assert.True(summary.Ok);
Assert.Equal(4, summary.Total);
Assert.Equal(4, summary.Migrated);
Assert.Equal(4, provisioner.ProvisionedSchemaNames.Count);
// По одной записи на страницу: 0,1,2,3, затем пустая страница завершает обход.
Assert.Equal(5, repository.PageRequests.Count);
Assert.Equal(0, repository.PageRequests[0].Offset);
Assert.Equal(3, repository.PageRequests[3].Offset);
}
[Fact]
public async Task MigrateAllAsync_SchemaFailsWithinPage_ContinuesOthers()
{
TenantRecordDto[] records = Records(3);
var failingSchema = $"tenant_{records[0].Id:N}";
var provisioner = new FailingTenantProvisioner(failingSchema);
var repository = new FakeTenantRepository(records);
var service = NewService(repository, provisioner);
TenantMigrationSummary summary = await service.MigrateAllAsync(1, 2, CancellationToken.None);
Assert.False(summary.Ok);
Assert.Equal(3, summary.Total);
Assert.Equal(2, summary.Migrated);
Assert.Equal(1, summary.Failed);
Assert.Equal(new[] { failingSchema }, summary.FailedSchemas);
// Уцелевшие схемы обеих страниц провижинены, сбойная — нет.
Assert.Equal(2, provisioner.ProvisionedSchemaNames.Count);
}
[Fact]
public async Task MigrateAllAsync_NonPositiveParameters_ClampToMinimum()
{
var provisioner = new FakeTenantProvisioner();
var repository = new FakeTenantRepository(Records(3));
var service = NewService(repository, provisioner);
TenantMigrationSummary summary = await service.MigrateAllAsync(0, 0, CancellationToken.None);
Assert.True(summary.Ok);
Assert.Equal(3, summary.Total);
Assert.Equal(3, summary.Migrated);
// Кламп pageSize → 1: по одной записи на страницу (плюс пустая на завершение).
Assert.All(repository.PageRequests, request => Assert.Equal(1, request.Limit));
}
[Fact]
public async Task MigrateAllAsync_WithExplicitParallelism_ProvisionsEveryTenant()
{
var provisioner = new FakeTenantProvisioner();
var repository = new FakeTenantRepository(
TenantRecord(TenantA, "A"),
TenantRecord(TenantB, "B"));
var service = NewService(repository, provisioner);
TenantMigrationSummary summary = await service.MigrateAllAsync(2, CancellationToken.None);
Assert.True(summary.Ok);
Assert.Equal(2, summary.Total);
Assert.Equal(2, summary.Migrated);
}
// Создаёт сервис на фейках (логирование не проверяется).
// repository: Реестр тенантов.
// provisioner: Провижинер схем.
@@ -84,4 +176,19 @@ public sealed class TenantSchemaMigrationServiceTests
// Возвращает: Запись реестра.
private static TenantRecordDto TenantRecord(Guid id, string name)
=> new(id, name, TenantStatuses.Active, DateTimeOffset.UtcNow);
// Последовательность тенантов с детерминированными Id (для проверки шардирования).
// count: Сколько записей создать.
// Возвращает: Записи реестра в порядке обхода.
private static TenantRecordDto[] Records(int count)
{
var records = new TenantRecordDto[count];
for (int index = 0; index < count; index++)
{
var id = Guid.Parse($"33333333-3333-3333-3333-{index + 1:D12}");
records[index] = TenantRecord(id, $"T{index + 1}");
}
return records;
}
}
@@ -27,6 +27,16 @@ public sealed class FakeTenantRegistry : ITenantRepository
public Task<IReadOnlyList<TenantRecordDto>> ListAsync(CancellationToken ct) =>
Task.FromResult<IReadOnlyList<TenantRecordDto>>(_tenants.ToList());
/// <inheritdoc />
public Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct)
{
IReadOnlyList<TenantRecordDto> page = _tenants.Skip(offset).Take(limit).ToList();
return Task.FromResult(page);
}
/// <inheritdoc />
public Task CreateAsync(TenantRecordDto tenant, CancellationToken ct) =>
throw new NotSupportedException("CreateAsync не используется тестами ингресса");
@@ -9,6 +9,7 @@ namespace Deal.Tests.Unit.Modules.Tenants;
public sealed class FakeTenantRepository : ITenantRepository
{
private readonly IReadOnlyList<TenantRecordDto> _tenants;
private readonly List<(int Offset, int Limit)> _pageRequests = [];
/// <summary>
/// Создаёт реестр с фиксированным списком тенантов.
@@ -19,10 +20,26 @@ public sealed class FakeTenantRepository : ITenantRepository
_tenants = tenants.ToList();
}
/// <summary>
/// Запрошенные страницы (offset, limit) в порядке обращений.
/// </summary>
public IReadOnlyList<(int Offset, int Limit)> PageRequests => _pageRequests;
/// <inheritdoc />
public Task<IReadOnlyList<TenantRecordDto>> ListAsync(CancellationToken ct) =>
Task.FromResult<IReadOnlyList<TenantRecordDto>>(_tenants.ToList());
/// <inheritdoc />
public Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct)
{
_pageRequests.Add((offset, limit));
IReadOnlyList<TenantRecordDto> page = _tenants.Skip(offset).Take(limit).ToList();
return Task.FromResult(page);
}
/// <inheritdoc />
public Task<TenantRecordDto?> FindByIdAsync(Guid id, CancellationToken ct) =>
throw new NotSupportedException("FindByIdAsync не используется тестами реестра");
@@ -39,6 +39,16 @@ public sealed class FakeTenantStore : ITenantRepository
public Task<IReadOnlyList<TenantRecordDto>> ListAsync(CancellationToken ct) =>
Task.FromResult<IReadOnlyList<TenantRecordDto>>(_tenants.ToList());
/// <inheritdoc />
public Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct)
{
IReadOnlyList<TenantRecordDto> page = _tenants.Skip(offset).Take(limit).ToList();
return Task.FromResult(page);
}
/// <inheritdoc />
public Task<bool> UpdateStatusAsync(
Guid id,
@@ -460,6 +460,13 @@ public sealed class StorageTickSchedulerTests
public Task<IReadOnlyList<TenantRecordDto>> ListAsync(CancellationToken ct) =>
throw new InvalidOperationException("реестр тенантов недоступен");
/// <inheritdoc />
public Task<IReadOnlyList<TenantRecordDto>> ListPageAsync(
int offset,
int limit,
CancellationToken ct) =>
throw new InvalidOperationException("реестр тенантов недоступен");
/// <inheritdoc />
public Task<TenantRecordDto?> FindByIdAsync(Guid id, CancellationToken ct) =>
throw new NotSupportedException();