Files
Deal/src/telegram-service/Deal.Telegram/Dialogs/BackfillService.cs
T
Rustam Khalimov e3a2692507 Добить структуру Api, Contracts, SharedKernel и сервисов
Deal.Api/Http -> Services/Models/Extensions; Contracts/Integrations
и SharedKernel/Tenants -> Abstractions/Models; extension-классы
telegram/ml -> Extensions. namespace/using/FQN мигрированы, using
дедуплицированы.
2026-09-11 13:25:18 +03:00

207 lines
10 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
using System.Collections.Concurrent;
using Deal.Grpc.Telegram;
using Deal.Telegram.Core;
using Deal.Telegram.Sessions;
using Deal.Telegram.Telegram;
using Deal.Telegram.Extensions;
namespace Deal.Telegram.Dialogs;
/// <summary>
/// Backfill диалога: перечитывание последних ~10 сообщений и отправка их в ядро потоком PushMessage
/// (план Task 10, Dialogs/BackfillService.cs; прототип backfill_dialog/_backfill_dialogs L331390).
///
/// 1:1 с прототипом: сообщения читаются от старых к новым (reversed — как реальный поток), между
/// отправками — анти-бан-пауза 1.5–3 с/сообщение (Ruling 3, BACKFILL_PER_MESSAGE L35); между
/// диалогами одного тенанта — пауза 3–6 с (BACKFILL_PER_DIALOG L36, python sleep между диалогами
/// L347). В конце — read-ack (снять «новое» в Telegram). Ошибка середины потока прерывает backfill
/// без read-ack: сообщения, не дошедшие до ядра, останутся непрочитанными и будут догнаны
/// realtime_sweep; дубли уже доставленных в ядро не растут (дубль-гвард dialog+msgId, Ruling 7).
/// Параметр force («Перечитать» по кнопке) прототип использует против своего флага backfilled;
/// в разделении флаг живёт в БД ядра, поэтому RPC исполняется всегда — дубли гасит ядро.
/// </summary>
public sealed class BackfillService
{
/// <summary>
/// Сколько последних сообщений читает backfill (прототип: limit=10, L371).
/// </summary>
public const int MessagesLimit = 10;
/// <summary>
/// Нижняя граница анти-бан-паузы между сообщениями (BACKFILL_PER_MESSAGE, L35).
/// </summary>
public const double MinPerMessageDelaySeconds = 1.5;
/// <summary>
/// Верхняя граница анти-бан-паузы между сообщениями (BACKFILL_PER_MESSAGE, L35).
/// </summary>
public const double MaxPerMessageDelaySeconds = 3.0;
/// <summary>
/// Нижняя граница паузы между диалогами одного тенанта (BACKFILL_PER_DIALOG, L36).
/// </summary>
public const double MinPerDialogDelaySeconds = 3.0;
/// <summary>
/// Верхняя граница паузы между диалогами одного тенанта (BACKFILL_PER_DIALOG, L36).
/// </summary>
public const double MaxPerDialogDelaySeconds = 6.0;
private readonly SessionFarm _sessionFarm;
private readonly ICoreIngressClient _ingress;
private readonly IBackfillPacer _pacer;
private readonly ILogger<BackfillService> _logger;
// Гард повторного входа: (тенант, диалог) уже перечитывается (python `_backfilling` L99).
private readonly ConcurrentDictionary<(string TenantId, string DialogId), byte> _running = new();
// Сериализация backfill'ов одного тенанта (python: один asyncio-loop на все диалоги).
private readonly ConcurrentDictionary<string, SemaphoreSlim> _tenantGates = new(StringComparer.Ordinal);
// Момент завершения последнего backfill тенанта (пауза 3–6 с между диалогами).
private readonly Dictionary<string, DateTime> _tenantLastFinishUtc = new(StringComparer.Ordinal);
private readonly object _lastFinishGate = new();
/// <summary>
/// Создаёт службу backfill'а.
/// </summary>
/// <param name="sessionFarm">Пул сессий тенантов (read-операции диалогов).</param>
/// <param name="ingress">Канал в ядро (PushMessage).</param>
/// <param name="pacer">Анти-бан-паузы (реальный — случайные, тесты — фейк).</param>
/// <param name="logger">Логгер.</param>
public BackfillService(
SessionFarm sessionFarm,
ICoreIngressClient ingress,
IBackfillPacer pacer,
ILogger<BackfillService> logger)
{
_sessionFarm = sessionFarm;
_ingress = ingress;
_pacer = pacer;
_logger = logger;
}
/// <summary>
/// Перечитывает последние сообщения диалога в ядро (Backfill RPC; кнопка «Перечитать»).
/// </summary>
/// <param name="tenantId">Id тенанта.</param>
/// <param name="dialogId">Подписанный id диалога.</param>
/// <param name="force">True — «Перечитать» по кнопке (семантика флага — в ядре, см. класс).</param>
/// <param name="cancellationToken">Отмена операции.</param>
/// <returns>Сколько сообщений отправлено в ядро (0 — диалог уже перечитывается/нет текстов).</returns>
public async Task<int> ExecuteAsync(
string tenantId,
string dialogId,
bool force,
CancellationToken cancellationToken)
{
var key = (tenantId, dialogId);
if (!_running.TryAdd(key, 0))
{
// Прототип L355–356: повторный вход в уже перечитываемый диалог → 0.
_logger.LogInformation("Аудит: backfill {TenantId} {DialogId} → пропущен (уже перечитывается)", tenantId, dialogId);
return 0;
}
try
{
SemaphoreSlim gate = _tenantGates.GetOrAdd(tenantId, static _ => new SemaphoreSlim(1, 1));
await gate.WaitAsync(cancellationToken).ConfigureAwait(false);
try
{
await EnsureDialogSpacingAsync(tenantId, cancellationToken).ConfigureAwait(false);
int processed = await BackfillDialogAsync(tenantId, dialogId, cancellationToken).ConfigureAwait(false);
lock (_lastFinishGate)
{
_tenantLastFinishUtc[tenantId] = DateTime.UtcNow;
}
_logger.LogInformation(
"Аудит: backfill {TenantId} {DialogId} → {Processed} сообщений в ядро (force {Force})",
tenantId,
dialogId,
processed,
force);
return processed;
}
finally
{
gate.Release();
}
}
finally
{
_running.TryRemove(key, out _);
}
}
// Пауза 36 с между backfill'ами разных диалогов одного тенанта (python L347).
// tenantId: Id тенанта.
// cancellationToken: Отмена операции.
private async Task EnsureDialogSpacingAsync(string tenantId, CancellationToken cancellationToken)
{
DateTime? lastFinishUtc;
lock (_lastFinishGate)
{
lastFinishUtc = _tenantLastFinishUtc.TryGetValue(tenantId, out DateTime value) ? value : null;
}
if (lastFinishUtc is null)
{
return;
}
double elapsedSeconds = (DateTime.UtcNow - lastFinishUtc.Value).TotalSeconds;
if (elapsedSeconds >= MaxPerDialogDelaySeconds)
{
return;
}
double minSeconds = Math.Max(0, MinPerDialogDelaySeconds - elapsedSeconds);
double maxSeconds = MaxPerDialogDelaySeconds - elapsedSeconds;
if (maxSeconds < minSeconds)
{
maxSeconds = minSeconds;
}
await _pacer.WaitAsync(minSeconds, maxSeconds, cancellationToken).ConfigureAwait(false);
}
// Читает и отправляет сообщения одного диалога (паузы между сообщениями, read-ack).
// tenantId: Id тенанта.
// dialogId: Подписанный id диалога.
// cancellationToken: Отмена операции.
// Возвращает: Сколько сообщений отправлено в ядро.
private async Task<int> BackfillDialogAsync(
string tenantId,
string dialogId,
CancellationToken cancellationToken)
{
IReadOnlyList<TelegramMessage> messages = await _sessionFarm
.GetMessagesAsync(tenantId, dialogId, MessagesLimit, cancellationToken)
.ConfigureAwait(false);
int processed = 0;
foreach (TelegramMessage message in messages.Reverse())
{
// От старых к новым — как реальный поток (прототип L372: `for m in reversed(msgs)`).
PushMessageReply reply = await _ingress
.PushMessageAsync(tenantId, DialogProtoMapper.ToPushRequest(message), cancellationToken)
.ConfigureAwait(false);
if (reply.Accepted && !reply.Duplicate)
{
processed++;
}
// Анти-бан между сообщениями (прототип L381; Ruling 3).
await _pacer.WaitAsync(MinPerMessageDelaySeconds, MaxPerMessageDelaySeconds, cancellationToken).ConfigureAwait(false);
}
// «Перечитали» — снимаем «новое» в Telegram (прототип L383–386). При ошибке выше read-ack
// не выполняется: неотправленное останется непрочитанным и догонится realtime_sweep.
await _sessionFarm.MarkReadAsync(tenantId, dialogId, cancellationToken).ConfigureAwait(false);
return processed;
}
}