Убрать неиспользуемые using по код-стайлу
Прогон dotnet format (IDE0005) по 4 решениям: удалены лишние using, оставшиеся после миграции namespace (676 файлов).
This commit is contained in:
@@ -1,436 +1,427 @@
|
||||
using System.Text.Json;
|
||||
using Deal.Api.Events;
|
||||
using Deal.Contracts.Integrations.Models;
|
||||
using Deal.Grpc.Telegram;
|
||||
using Deal.Modules.Pipeline.Application.Abstractions;
|
||||
using Deal.Modules.Pipeline.Application.Models;
|
||||
using Deal.Modules.Pipeline.Application.Registrars;
|
||||
using Deal.Modules.Pipeline.Application.Services;
|
||||
using Deal.Modules.Settings.Application.Abstractions;
|
||||
using Deal.Modules.Settings.Application.Models;
|
||||
using Deal.Modules.Settings.Application.Registrars;
|
||||
using Deal.Modules.Settings.Application.Services;
|
||||
using Deal.Modules.Telegram.Application;
|
||||
using Deal.Modules.Tenants.Application.Abstractions;
|
||||
using Deal.Modules.Tenants.Application.Extensions;
|
||||
using Deal.Modules.Tenants.Application.Models;
|
||||
using Deal.Modules.Tenants.Application.Registrars;
|
||||
using Deal.Modules.Tenants.Application.Services;
|
||||
using Deal.SharedKernel.Tenants.Abstractions;
|
||||
using Deal.SharedKernel.Tenants.Models;
|
||||
using Grpc.Core;
|
||||
using Deal.Api.Services;
|
||||
using Deal.Api.Dtos;
|
||||
|
||||
namespace Deal.Api.Telegram;
|
||||
|
||||
/// <summary>
|
||||
/// gRPC-сервер входящего потока telegram-service → ядро (план Task 12, L361–377; Ruling 1/7).
|
||||
///
|
||||
/// Реализация серверной стороны Deal.Grpc.Telegram.IngressService (telegram.proto, L380–395):
|
||||
/// PushMessage — новое/догоняющее сообщение мониторящегося диалога в очередь пайплайна
|
||||
/// (<see cref="PipelineIngestService.EnqueueAsync"/>, тот же контракт, что приём сообщений пайплайна) в схеме тенанта
|
||||
/// + превью (DialogsService.SavePreview: TgMessages + «последнее сообщение» каталога, Ruling 7);
|
||||
/// SyncDialogs — применение каталога диалогов (DialogsService.SyncFromTelegram) и ответ со списком
|
||||
/// monitored id (зеркало сервиса); ReportStatus — статус аккаунта в KV (tgStatus/tgAccount) + SSE
|
||||
/// system_status/тосты на переходах фаз.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// Tenant-id берётся ТОЛЬКО из gRPC-metadata (полю в теле не доверяем — Ruling 1), принадлежность
|
||||
/// подтверждается реестром тенантов (public.tenants), затем для работы открывается собственный scope
|
||||
/// с <c>ITenantContext.SetTenant</c> (эталон PipelineWorkerScheduler, L169–213): tenant-scoped адаптеры
|
||||
/// (PipelineStore/SettingsStore) строятся от схемы тенанта. Неизвестный тенант/сбой схемы — RPC не падает:
|
||||
/// ответ не-принято (accepted=false / ok=false, план Task 12) + лог аудита (Ruling 13); недоступный сервис
|
||||
/// догоняет упущенное realtime-sweep (контракт README).
|
||||
/// <para>
|
||||
/// Полная синхронизация каталога (применение entries к таблице Dialogs, ответ = список monitored id) —
|
||||
/// модуль Deal.Modules.Telegram (план Task 13): DialogsService.SyncFromTelegram (upsert/удаление, авто-
|
||||
/// мониторинг новых по autoMonitorNew), превью сообщений — DialogsService.SavePreview (PushMessage).
|
||||
/// </para>
|
||||
/// </remarks>
|
||||
public sealed class TelegramIngressService(
|
||||
IServiceScopeFactory scopeFactory,
|
||||
SseBroker broker,
|
||||
ILogger<TelegramIngressService> logger) : IngressService.IngressServiceBase
|
||||
{
|
||||
/// <summary>
|
||||
/// Ключ gRPC-metadata с id тенанта (единственный источник принадлежности — Ruling 1).
|
||||
/// </summary>
|
||||
public const string TenantIdMetadataKey = "tenant-id";
|
||||
|
||||
// Тип SSE-события статуса Telegram (фронт по нему перечитывает GET /api/tg/status, Ruling 7).
|
||||
private const string SystemStatusEventType = "system_status";
|
||||
|
||||
// Тип SSE-события тоста (Ruling 5; api.js L79 слушает 'toast').
|
||||
private const string ToastEventType = "toast";
|
||||
|
||||
// Текст тоста подключения (Ruling 7, 1:1 с прототипом).
|
||||
private const string ConnectedToastText = "Telegram подключён, сессия сохранена";
|
||||
|
||||
// Текст тоста отключения (Ruling 7, 1:1 с прототипом).
|
||||
private const string DisconnectedToastText = "Telegram отключён";
|
||||
|
||||
// Иконка тоста подключения (из набора Icon.vue фронта).
|
||||
private const string ConnectedToastIcon = "send";
|
||||
|
||||
// Иконка тоста отключения (из набора Icon.vue фронта).
|
||||
private const string DisconnectedToastIcon = "logout";
|
||||
|
||||
// Деталь отказа: metadata tenant-id отсутствует (UNAUTHENTICATED, README src/contracts).
|
||||
private const string MissingTenantIdDetail = "tenant-id отсутствует в metadata";
|
||||
|
||||
// Опции JSON KV-статуса: camelCase (1:1 с wire-именами) + терпимость регистра при чтении.
|
||||
private static readonly JsonSerializerOptions StatusJsonOptions = new()
|
||||
{
|
||||
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
|
||||
PropertyNameCaseInsensitive = true,
|
||||
};
|
||||
|
||||
/// <summary>
|
||||
/// PushMessage — сообщение диалога в очередь пайплайна тенанта + превью (PushMessageRequest, Ruling 7).
|
||||
/// </summary>
|
||||
/// <remarks>Дубль dialog_id+msg_id уже в очереди — duplicate=true, очередь не растёт (гвард
|
||||
/// PipelineIngestService). Пустой текст/диалог — no-op приёма (accepted=false, контракт proto).
|
||||
/// После постановки в очередь пишется превью (DialogsService.SavePreview: строка TgMessages
|
||||
/// «m_<dialog>_<msg>» + «последнее сообщение» каталога — 1:1 _on_message python L270–274);
|
||||
/// сбой превью не влияет на приём (accepted определён очередью, лог дебага).
|
||||
/// Неизвестный тенант или сбой схемы/БД — не-принято (accepted=false) без исключения RPC.</remarks>
|
||||
/// <param name="request">Сообщение из потока telegram-service.</param>
|
||||
/// <param name="context">Контекст вызова (metadata tenant-id + service-token).</param>
|
||||
/// <returns>accepted — сообщение принято (либо дубль), duplicate — уже было в очереди.</returns>
|
||||
public override async Task<PushMessageReply> PushMessage(PushMessageRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new PushMessageReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
// Resolve ПОСЛЕ SetTenant: TenantDbContext (и его адаптеры) строятся от схемы текущего тенанта.
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
PipelineIngestService ingest = tenantScope.ServiceProvider.GetRequiredService<PipelineIngestService>();
|
||||
|
||||
PipelineIngestResultDto result = await ingest.EnqueueAsync(
|
||||
new QueuedMessage
|
||||
{
|
||||
DialogId = request.DialogId,
|
||||
ChannelName = request.ChannelName,
|
||||
ChannelHandle = request.ChannelHandle,
|
||||
ChannelHue = request.ChannelHue,
|
||||
MsgId = request.HasMsgId ? request.MsgId : null,
|
||||
Text = request.Text,
|
||||
MsgAtMs = request.HasMsgAt ? request.MsgAt : null,
|
||||
},
|
||||
context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
await SavePreviewSafelyAsync(tenantScope, tenant, request, context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: PushMessage {TenantId} диалог {DialogId} msg {MsgId} → {Outcome}",
|
||||
tenant.Id,
|
||||
request.DialogId,
|
||||
request.HasMsgId ? request.MsgId.ToString() : "-",
|
||||
result.Duplicate ? "duplicate" : result.Id is null ? "no-op" : "queued");
|
||||
|
||||
return new PushMessageReply
|
||||
{
|
||||
Accepted = result.Id is not null || result.Duplicate,
|
||||
Duplicate = result.Duplicate,
|
||||
};
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
// Сбой схемы/БД тенанта (напр. схема ещё не провижинена): RPC не падает — reply not-accepted
|
||||
// (план Task 12), упущенное сообщение при необходимости догонит realtime-sweep сервиса.
|
||||
logger.LogWarning(exception, "Аудит: PushMessage {TenantId} → не принято (сбой схемы/БД)", tenant.Id);
|
||||
return new PushMessageReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// SyncDialogs — синхронизация каталога диалогов аккаунта (Ruling 7, L386–390).
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// Модуль Deal.Modules.Telegram (план Task 13) применяет entries к таблице Dialogs
|
||||
/// (DialogsService.SyncFromTelegram: upsert + удаление отсутствующих; авто-мониторинг новых — по
|
||||
/// настройке autoMonitorNew). Ответ несёт актуальный список monitored id — по нему telegram-service
|
||||
/// держит своё зеркало мониторинга в памяти (обновляется ответом SyncDialogs и командой SetMonitor,
|
||||
/// Ruling 7) и фильтрует события realtime.
|
||||
/// </remarks>
|
||||
/// <param name="request">Актуальный каталог диалогов (entries).</param>
|
||||
/// <param name="context">Контекст вызова.</param>
|
||||
/// <returns>monitored_ids — диалоги с включённым мониторингом после применения каталога.</returns>
|
||||
public override async Task<SyncDialogsReply> SyncDialogs(SyncDialogsRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new SyncDialogsReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
DialogsService dialogs = tenantScope.ServiceProvider.GetRequiredService<DialogsService>();
|
||||
|
||||
List<TelegramDialogEntryDto> entries = new(request.Entries.Count);
|
||||
foreach (DialogEntry entry in request.Entries)
|
||||
{
|
||||
entries.Add(new TelegramDialogEntryDto(entry.Id, entry.Name, entry.Username, entry.Kind, entry.Hue));
|
||||
}
|
||||
|
||||
int synced = await dialogs.SyncFromTelegramAsync(entries, context.CancellationToken).ConfigureAwait(false);
|
||||
IReadOnlyCollection<string> monitoredIds =
|
||||
await dialogs.ListMonitoredIdsAsync(context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: SyncDialogs {TenantId}: каталог {Count} → применено {Synced}, monitored {Monitored}",
|
||||
tenant.Id,
|
||||
request.Entries.Count,
|
||||
synced,
|
||||
monitoredIds.Count);
|
||||
|
||||
var reply = new SyncDialogsReply();
|
||||
reply.MonitoredIds.AddRange(monitoredIds);
|
||||
return reply;
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: SyncDialogs {TenantId} → каталог не применён (сбой схемы/БД)", tenant.Id);
|
||||
return new SyncDialogsReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// ReportStatus — статус аккаунта в KV + SSE system_status/тосты на переходах фаз (Ruling 7).
|
||||
/// </summary>
|
||||
/// <remarks>KV tgStatus (снимок без account) и tgAccount (JSON-строка) пишутся в схему тенанта;
|
||||
/// system_status публикуется на каждый репорт (фронт перечитывает /api/tg/status), тосты — только на
|
||||
/// переходы connected: false→true «Telegram подключён, сессия сохранена», true→false «Telegram отключён»
|
||||
/// (сервис шлёт статус по событию и heartbeat'ом — без гарда переходов тосты дублировались бы).
|
||||
/// Неизвестный тенант/сбой схемы — ok=false без исключения RPC (план Task 12).</remarks>
|
||||
/// <param name="request">Статус аккаунта из _publish_status прототипа.</param>
|
||||
/// <param name="context">Контекст вызова.</param>
|
||||
/// <returns>ok — статус принят и сохранён.</returns>
|
||||
public override async Task<ReportStatusReply> ReportStatus(ReportStatusRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new ReportStatusReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
ISettingsStore settings = tenantScope.ServiceProvider.GetRequiredService<ISettingsStore>();
|
||||
|
||||
TgReportedStatus current = ToReportedStatus(request);
|
||||
TgReportedStatus? previous = await ReadPreviousStatusAsync(settings, context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
// SSE до записи KV: канал тенанта обновляется и при сбое записи (следующий репорт перепишет KV).
|
||||
PublishStatusEvents(tenant.Id, previous, current);
|
||||
|
||||
await settings.SetAsync(SettingsKeys.TgStatus, ToJson(current), context.CancellationToken).ConfigureAwait(false);
|
||||
await settings.SetAsync(SettingsKeys.TgAccount, JsonSerializer.Serialize(request.Account), context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: ReportStatus {TenantId} → фаза {Phase}, connected {Connected}",
|
||||
tenant.Id,
|
||||
request.Phase,
|
||||
request.Connected);
|
||||
return new ReportStatusReply { Ok = true };
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: ReportStatus {TenantId} → не сохранён (сбой схемы/БД)", tenant.Id);
|
||||
return new ReportStatusReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
// Разрешает тенанта запроса: metadata tenant-id → реестр public.tenants.
|
||||
// Отсутствующий/пустой tenant-id — RPC-отказ UNAUTHENTICATED (README: tenant-id обязателен).
|
||||
// Id не Guid либо записи нет в реестре — неизвестный тенант: лог аудита и null (RPC отвечает не-принято,
|
||||
// план Task 12: «для несуществующего тенанта не падает»).
|
||||
// context: Контекст вызова.
|
||||
// Возвращает: Запись тенанта реестра либо null (тенант неизвестен).
|
||||
private async Task<TenantRecordDto?> ResolveTenantAsync(ServerCallContext context)
|
||||
{
|
||||
string tenantId = RequireTenantIdMetadata(context);
|
||||
if (!Guid.TryParse(tenantId, out Guid tenantGuid))
|
||||
{
|
||||
logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (id не Guid)", context.Method, tenantId);
|
||||
return null;
|
||||
}
|
||||
|
||||
await using AsyncServiceScope registryScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantRepository repository = registryScope.ServiceProvider.GetRequiredService<ITenantRepository>();
|
||||
TenantRecordDto? tenant = await repository.FindByIdAsync(tenantGuid, context.CancellationToken).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (нет в реестре)", context.Method, tenantId);
|
||||
}
|
||||
|
||||
return tenant;
|
||||
}
|
||||
|
||||
// Читает tenant-id из metadata (обязателен; отсутствие — UNAUTHENTICATED, README).
|
||||
// context: Контекст вызова.
|
||||
// Возвращает: Значение tenant-id.
|
||||
private static string RequireTenantIdMetadata(ServerCallContext context)
|
||||
{
|
||||
string? tenantId = context.RequestHeaders.GetValue(TenantIdMetadataKey);
|
||||
if (string.IsNullOrWhiteSpace(tenantId))
|
||||
{
|
||||
throw new RpcException(new Status(StatusCode.Unauthenticated, MissingTenantIdDetail));
|
||||
}
|
||||
|
||||
return tenantId;
|
||||
}
|
||||
|
||||
// Пишет превью принятого сообщения (TgMessages + «последнее сообщение» каталога) без влияния на приём.
|
||||
// Ruling 7: PushMessage → EnqueueAsync + превью. Сбой превью (нет таблиц/строки каталога и т.п.)
|
||||
// не роняет RPC и не меняет accepted — очередь уже записана, упущенное догонит realtime-sweep (как
|
||||
// python: обновление last_text после enqueue в том же обработчике, ошибка не отменяет приём).
|
||||
// tenantScope: Scope тенанта (TenantDbContext построен на схеме тенанта).
|
||||
// tenant: Тенант канала (для лога аудита).
|
||||
// request: Сообщение PushMessage.
|
||||
// ct: Токен отмены.
|
||||
private async Task SavePreviewSafelyAsync(
|
||||
AsyncServiceScope tenantScope,
|
||||
TenantRecordDto tenant,
|
||||
PushMessageRequest request,
|
||||
CancellationToken ct)
|
||||
{
|
||||
try
|
||||
{
|
||||
DialogsService dialogs = tenantScope.ServiceProvider.GetRequiredService<DialogsService>();
|
||||
DateTimeOffset? msgAt = request.HasMsgAt ? DateTimeOffset.FromUnixTimeMilliseconds(request.MsgAt) : null;
|
||||
await dialogs.SavePreviewAsync(
|
||||
request.DialogId,
|
||||
request.HasMsgId ? request.MsgId : null,
|
||||
request.Text,
|
||||
msgAt,
|
||||
ct).ConfigureAwait(false);
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
// Превью — вторичная запись: приём не затронут (лог дебага, не ошибка RPC).
|
||||
logger.LogDebug(exception, "PushMessage {TenantId}: превью не сохранено (приём не затронут)", tenant.Id);
|
||||
}
|
||||
}
|
||||
|
||||
// Публикует SSE system_status (каждый репорт) и тосты на переходах connected.
|
||||
// tenantId: Тенант канала (реестровый Guid).
|
||||
// previous: Предыдущий снимок из KV (null — первый репорт).
|
||||
// current: Текущий снимок репорта.
|
||||
private void PublishStatusEvents(
|
||||
Guid tenantId,
|
||||
TgReportedStatus? previous,
|
||||
TgReportedStatus current)
|
||||
{
|
||||
broker.Publish(tenantId, SystemStatusEventType, current);
|
||||
if (previous is null)
|
||||
{
|
||||
// Первый репорт после старта сервиса: переходов нет, статус фронт получит по system_status.
|
||||
return;
|
||||
}
|
||||
|
||||
if (!previous.Connected && current.Connected)
|
||||
{
|
||||
PublishToast(tenantId, ConnectedToastText, ConnectedToastIcon);
|
||||
}
|
||||
else if (previous.Connected && !current.Connected)
|
||||
{
|
||||
PublishToast(tenantId, DisconnectedToastText, DisconnectedToastIcon);
|
||||
}
|
||||
}
|
||||
|
||||
// Публикует SSE-тост в канал тенанта (без подписчиков — no-op, Ruling 5).
|
||||
// tenantId: Тенант-получатель.
|
||||
// text: Текст тоста.
|
||||
// icon: Иконка тоста (набор Icon.vue фронта).
|
||||
private void PublishToast(
|
||||
Guid tenantId,
|
||||
string text,
|
||||
string icon)
|
||||
{
|
||||
broker.Publish(tenantId, ToastEventType, new { text, icon });
|
||||
}
|
||||
|
||||
// Снимок предыдущего статуса из KV tgStatus (нет записи/битый JSON — null).
|
||||
// settings: KV-хранилище настроек схемы тенанта.
|
||||
// ct: Токен отмены.
|
||||
// Возвращает: Предыдущий снимок либо null.
|
||||
private async Task<TgReportedStatus?> ReadPreviousStatusAsync(ISettingsStore settings, CancellationToken ct)
|
||||
{
|
||||
SettingValue? stored = await settings.GetAsync(SettingsKeys.TgStatus, ct).ConfigureAwait(false);
|
||||
if (stored is null)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
return JsonSerializer.Deserialize<TgReportedStatus>(stored.ValueJson, StatusJsonOptions);
|
||||
}
|
||||
catch (JsonException exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: KV tgStatus повреждён — переходы фаз не определяются");
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
// Маппит запрос ReportStatus в снимок KV (account живёт отдельным ключом tgAccount).
|
||||
// request: Запрос ReportStatus.
|
||||
// Возвращает: Снимок статуса.
|
||||
private static TgReportedStatus ToReportedStatus(ReportStatusRequest request) => new()
|
||||
{
|
||||
Phase = request.Phase,
|
||||
Connected = request.Connected,
|
||||
Listener = request.Listener,
|
||||
Error = request.HasError ? request.Error : null,
|
||||
QrUrl = request.HasQrUrl ? request.QrUrl : null,
|
||||
};
|
||||
|
||||
// Сериализует снимок в JSON (camelCase, конвенция value_json).
|
||||
// status: Снимок статуса.
|
||||
// Возвращает: JSON-строка.
|
||||
private static string ToJson(TgReportedStatus status) => JsonSerializer.Serialize(status, StatusJsonOptions);
|
||||
}
|
||||
using System.Text.Json;
|
||||
using Deal.Api.Events;
|
||||
using Deal.Contracts.Integrations.Models;
|
||||
using Deal.Grpc.Telegram;
|
||||
using Deal.Modules.Pipeline.Application.Models;
|
||||
using Deal.Modules.Pipeline.Application.Services;
|
||||
using Deal.Modules.Settings.Application.Abstractions;
|
||||
using Deal.Modules.Settings.Application.Models;
|
||||
using Deal.Modules.Telegram.Application;
|
||||
using Deal.Modules.Tenants.Application.Abstractions;
|
||||
using Deal.Modules.Tenants.Application.Models;
|
||||
using Deal.SharedKernel.Tenants.Abstractions;
|
||||
using Deal.SharedKernel.Tenants.Models;
|
||||
using Grpc.Core;
|
||||
|
||||
namespace Deal.Api.Telegram;
|
||||
|
||||
/// <summary>
|
||||
/// gRPC-сервер входящего потока telegram-service → ядро (план Task 12, L361–377; Ruling 1/7).
|
||||
///
|
||||
/// Реализация серверной стороны Deal.Grpc.Telegram.IngressService (telegram.proto, L380–395):
|
||||
/// PushMessage — новое/догоняющее сообщение мониторящегося диалога в очередь пайплайна
|
||||
/// (<see cref="PipelineIngestService.EnqueueAsync"/>, тот же контракт, что приём сообщений пайплайна) в схеме тенанта
|
||||
/// + превью (DialogsService.SavePreview: TgMessages + «последнее сообщение» каталога, Ruling 7);
|
||||
/// SyncDialogs — применение каталога диалогов (DialogsService.SyncFromTelegram) и ответ со списком
|
||||
/// monitored id (зеркало сервиса); ReportStatus — статус аккаунта в KV (tgStatus/tgAccount) + SSE
|
||||
/// system_status/тосты на переходах фаз.
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// Tenant-id берётся ТОЛЬКО из gRPC-metadata (полю в теле не доверяем — Ruling 1), принадлежность
|
||||
/// подтверждается реестром тенантов (public.tenants), затем для работы открывается собственный scope
|
||||
/// с <c>ITenantContext.SetTenant</c> (эталон PipelineWorkerScheduler, L169–213): tenant-scoped адаптеры
|
||||
/// (PipelineStore/SettingsStore) строятся от схемы тенанта. Неизвестный тенант/сбой схемы — RPC не падает:
|
||||
/// ответ не-принято (accepted=false / ok=false, план Task 12) + лог аудита (Ruling 13); недоступный сервис
|
||||
/// догоняет упущенное realtime-sweep (контракт README).
|
||||
/// <para>
|
||||
/// Полная синхронизация каталога (применение entries к таблице Dialogs, ответ = список monitored id) —
|
||||
/// модуль Deal.Modules.Telegram (план Task 13): DialogsService.SyncFromTelegram (upsert/удаление, авто-
|
||||
/// мониторинг новых по autoMonitorNew), превью сообщений — DialogsService.SavePreview (PushMessage).
|
||||
/// </para>
|
||||
/// </remarks>
|
||||
public sealed class TelegramIngressService(
|
||||
IServiceScopeFactory scopeFactory,
|
||||
SseBroker broker,
|
||||
ILogger<TelegramIngressService> logger) : IngressService.IngressServiceBase
|
||||
{
|
||||
/// <summary>
|
||||
/// Ключ gRPC-metadata с id тенанта (единственный источник принадлежности — Ruling 1).
|
||||
/// </summary>
|
||||
public const string TenantIdMetadataKey = "tenant-id";
|
||||
|
||||
// Тип SSE-события статуса Telegram (фронт по нему перечитывает GET /api/tg/status, Ruling 7).
|
||||
private const string SystemStatusEventType = "system_status";
|
||||
|
||||
// Тип SSE-события тоста (Ruling 5; api.js L79 слушает 'toast').
|
||||
private const string ToastEventType = "toast";
|
||||
|
||||
// Текст тоста подключения (Ruling 7, 1:1 с прототипом).
|
||||
private const string ConnectedToastText = "Telegram подключён, сессия сохранена";
|
||||
|
||||
// Текст тоста отключения (Ruling 7, 1:1 с прототипом).
|
||||
private const string DisconnectedToastText = "Telegram отключён";
|
||||
|
||||
// Иконка тоста подключения (из набора Icon.vue фронта).
|
||||
private const string ConnectedToastIcon = "send";
|
||||
|
||||
// Иконка тоста отключения (из набора Icon.vue фронта).
|
||||
private const string DisconnectedToastIcon = "logout";
|
||||
|
||||
// Деталь отказа: metadata tenant-id отсутствует (UNAUTHENTICATED, README src/contracts).
|
||||
private const string MissingTenantIdDetail = "tenant-id отсутствует в metadata";
|
||||
|
||||
// Опции JSON KV-статуса: camelCase (1:1 с wire-именами) + терпимость регистра при чтении.
|
||||
private static readonly JsonSerializerOptions StatusJsonOptions = new()
|
||||
{
|
||||
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
|
||||
PropertyNameCaseInsensitive = true,
|
||||
};
|
||||
|
||||
/// <summary>
|
||||
/// PushMessage — сообщение диалога в очередь пайплайна тенанта + превью (PushMessageRequest, Ruling 7).
|
||||
/// </summary>
|
||||
/// <remarks>Дубль dialog_id+msg_id уже в очереди — duplicate=true, очередь не растёт (гвард
|
||||
/// PipelineIngestService). Пустой текст/диалог — no-op приёма (accepted=false, контракт proto).
|
||||
/// После постановки в очередь пишется превью (DialogsService.SavePreview: строка TgMessages
|
||||
/// «m_<dialog>_<msg>» + «последнее сообщение» каталога — 1:1 _on_message python L270–274);
|
||||
/// сбой превью не влияет на приём (accepted определён очередью, лог дебага).
|
||||
/// Неизвестный тенант или сбой схемы/БД — не-принято (accepted=false) без исключения RPC.</remarks>
|
||||
/// <param name="request">Сообщение из потока telegram-service.</param>
|
||||
/// <param name="context">Контекст вызова (metadata tenant-id + service-token).</param>
|
||||
/// <returns>accepted — сообщение принято (либо дубль), duplicate — уже было в очереди.</returns>
|
||||
public override async Task<PushMessageReply> PushMessage(PushMessageRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new PushMessageReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
// Resolve ПОСЛЕ SetTenant: TenantDbContext (и его адаптеры) строятся от схемы текущего тенанта.
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
PipelineIngestService ingest = tenantScope.ServiceProvider.GetRequiredService<PipelineIngestService>();
|
||||
|
||||
PipelineIngestResultDto result = await ingest.EnqueueAsync(
|
||||
new QueuedMessage
|
||||
{
|
||||
DialogId = request.DialogId,
|
||||
ChannelName = request.ChannelName,
|
||||
ChannelHandle = request.ChannelHandle,
|
||||
ChannelHue = request.ChannelHue,
|
||||
MsgId = request.HasMsgId ? request.MsgId : null,
|
||||
Text = request.Text,
|
||||
MsgAtMs = request.HasMsgAt ? request.MsgAt : null,
|
||||
},
|
||||
context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
await SavePreviewSafelyAsync(tenantScope, tenant, request, context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: PushMessage {TenantId} диалог {DialogId} msg {MsgId} → {Outcome}",
|
||||
tenant.Id,
|
||||
request.DialogId,
|
||||
request.HasMsgId ? request.MsgId.ToString() : "-",
|
||||
result.Duplicate ? "duplicate" : result.Id is null ? "no-op" : "queued");
|
||||
|
||||
return new PushMessageReply
|
||||
{
|
||||
Accepted = result.Id is not null || result.Duplicate,
|
||||
Duplicate = result.Duplicate,
|
||||
};
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
// Сбой схемы/БД тенанта (напр. схема ещё не провижинена): RPC не падает — reply not-accepted
|
||||
// (план Task 12), упущенное сообщение при необходимости догонит realtime-sweep сервиса.
|
||||
logger.LogWarning(exception, "Аудит: PushMessage {TenantId} → не принято (сбой схемы/БД)", tenant.Id);
|
||||
return new PushMessageReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// SyncDialogs — синхронизация каталога диалогов аккаунта (Ruling 7, L386–390).
|
||||
/// </summary>
|
||||
/// <remarks>
|
||||
/// Модуль Deal.Modules.Telegram (план Task 13) применяет entries к таблице Dialogs
|
||||
/// (DialogsService.SyncFromTelegram: upsert + удаление отсутствующих; авто-мониторинг новых — по
|
||||
/// настройке autoMonitorNew). Ответ несёт актуальный список monitored id — по нему telegram-service
|
||||
/// держит своё зеркало мониторинга в памяти (обновляется ответом SyncDialogs и командой SetMonitor,
|
||||
/// Ruling 7) и фильтрует события realtime.
|
||||
/// </remarks>
|
||||
/// <param name="request">Актуальный каталог диалогов (entries).</param>
|
||||
/// <param name="context">Контекст вызова.</param>
|
||||
/// <returns>monitored_ids — диалоги с включённым мониторингом после применения каталога.</returns>
|
||||
public override async Task<SyncDialogsReply> SyncDialogs(SyncDialogsRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new SyncDialogsReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
DialogsService dialogs = tenantScope.ServiceProvider.GetRequiredService<DialogsService>();
|
||||
|
||||
List<TelegramDialogEntryDto> entries = new(request.Entries.Count);
|
||||
foreach (DialogEntry entry in request.Entries)
|
||||
{
|
||||
entries.Add(new TelegramDialogEntryDto(entry.Id, entry.Name, entry.Username, entry.Kind, entry.Hue));
|
||||
}
|
||||
|
||||
int synced = await dialogs.SyncFromTelegramAsync(entries, context.CancellationToken).ConfigureAwait(false);
|
||||
IReadOnlyCollection<string> monitoredIds =
|
||||
await dialogs.ListMonitoredIdsAsync(context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: SyncDialogs {TenantId}: каталог {Count} → применено {Synced}, monitored {Monitored}",
|
||||
tenant.Id,
|
||||
request.Entries.Count,
|
||||
synced,
|
||||
monitoredIds.Count);
|
||||
|
||||
var reply = new SyncDialogsReply();
|
||||
reply.MonitoredIds.AddRange(monitoredIds);
|
||||
return reply;
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: SyncDialogs {TenantId} → каталог не применён (сбой схемы/БД)", tenant.Id);
|
||||
return new SyncDialogsReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// ReportStatus — статус аккаунта в KV + SSE system_status/тосты на переходах фаз (Ruling 7).
|
||||
/// </summary>
|
||||
/// <remarks>KV tgStatus (снимок без account) и tgAccount (JSON-строка) пишутся в схему тенанта;
|
||||
/// system_status публикуется на каждый репорт (фронт перечитывает /api/tg/status), тосты — только на
|
||||
/// переходы connected: false→true «Telegram подключён, сессия сохранена», true→false «Telegram отключён»
|
||||
/// (сервис шлёт статус по событию и heartbeat'ом — без гарда переходов тосты дублировались бы).
|
||||
/// Неизвестный тенант/сбой схемы — ok=false без исключения RPC (план Task 12).</remarks>
|
||||
/// <param name="request">Статус аккаунта из _publish_status прототипа.</param>
|
||||
/// <param name="context">Контекст вызова.</param>
|
||||
/// <returns>ok — статус принят и сохранён.</returns>
|
||||
public override async Task<ReportStatusReply> ReportStatus(ReportStatusRequest request, ServerCallContext context)
|
||||
{
|
||||
TenantRecordDto? tenant = await ResolveTenantAsync(context).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
return new ReportStatusReply();
|
||||
}
|
||||
|
||||
await using AsyncServiceScope tenantScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantContext tenantContext = tenantScope.ServiceProvider.GetRequiredService<ITenantContext>();
|
||||
try
|
||||
{
|
||||
tenantContext.SetTenant(new TenantId(tenant.Id.ToString("N")));
|
||||
ISettingsStore settings = tenantScope.ServiceProvider.GetRequiredService<ISettingsStore>();
|
||||
|
||||
TgReportedStatus current = ToReportedStatus(request);
|
||||
TgReportedStatus? previous = await ReadPreviousStatusAsync(settings, context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
// SSE до записи KV: канал тенанта обновляется и при сбое записи (следующий репорт перепишет KV).
|
||||
PublishStatusEvents(tenant.Id, previous, current);
|
||||
|
||||
await settings.SetAsync(SettingsKeys.TgStatus, ToJson(current), context.CancellationToken).ConfigureAwait(false);
|
||||
await settings.SetAsync(SettingsKeys.TgAccount, JsonSerializer.Serialize(request.Account), context.CancellationToken).ConfigureAwait(false);
|
||||
|
||||
logger.LogInformation(
|
||||
"Аудит: ReportStatus {TenantId} → фаза {Phase}, connected {Connected}",
|
||||
tenant.Id,
|
||||
request.Phase,
|
||||
request.Connected);
|
||||
return new ReportStatusReply { Ok = true };
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: ReportStatus {TenantId} → не сохранён (сбой схемы/БД)", tenant.Id);
|
||||
return new ReportStatusReply();
|
||||
}
|
||||
finally
|
||||
{
|
||||
tenantContext.Reset();
|
||||
}
|
||||
}
|
||||
|
||||
// Разрешает тенанта запроса: metadata tenant-id → реестр public.tenants.
|
||||
// Отсутствующий/пустой tenant-id — RPC-отказ UNAUTHENTICATED (README: tenant-id обязателен).
|
||||
// Id не Guid либо записи нет в реестре — неизвестный тенант: лог аудита и null (RPC отвечает не-принято,
|
||||
// план Task 12: «для несуществующего тенанта не падает»).
|
||||
// context: Контекст вызова.
|
||||
// Возвращает: Запись тенанта реестра либо null (тенант неизвестен).
|
||||
private async Task<TenantRecordDto?> ResolveTenantAsync(ServerCallContext context)
|
||||
{
|
||||
string tenantId = RequireTenantIdMetadata(context);
|
||||
if (!Guid.TryParse(tenantId, out Guid tenantGuid))
|
||||
{
|
||||
logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (id не Guid)", context.Method, tenantId);
|
||||
return null;
|
||||
}
|
||||
|
||||
await using AsyncServiceScope registryScope = scopeFactory.CreateAsyncScope();
|
||||
ITenantRepository repository = registryScope.ServiceProvider.GetRequiredService<ITenantRepository>();
|
||||
TenantRecordDto? tenant = await repository.FindByIdAsync(tenantGuid, context.CancellationToken).ConfigureAwait(false);
|
||||
if (tenant is null)
|
||||
{
|
||||
logger.LogWarning("Аудит: ингресс {Action} → тенант {TenantId} неизвестен (нет в реестре)", context.Method, tenantId);
|
||||
}
|
||||
|
||||
return tenant;
|
||||
}
|
||||
|
||||
// Читает tenant-id из metadata (обязателен; отсутствие — UNAUTHENTICATED, README).
|
||||
// context: Контекст вызова.
|
||||
// Возвращает: Значение tenant-id.
|
||||
private static string RequireTenantIdMetadata(ServerCallContext context)
|
||||
{
|
||||
string? tenantId = context.RequestHeaders.GetValue(TenantIdMetadataKey);
|
||||
if (string.IsNullOrWhiteSpace(tenantId))
|
||||
{
|
||||
throw new RpcException(new Status(StatusCode.Unauthenticated, MissingTenantIdDetail));
|
||||
}
|
||||
|
||||
return tenantId;
|
||||
}
|
||||
|
||||
// Пишет превью принятого сообщения (TgMessages + «последнее сообщение» каталога) без влияния на приём.
|
||||
// Ruling 7: PushMessage → EnqueueAsync + превью. Сбой превью (нет таблиц/строки каталога и т.п.)
|
||||
// не роняет RPC и не меняет accepted — очередь уже записана, упущенное догонит realtime-sweep (как
|
||||
// python: обновление last_text после enqueue в том же обработчике, ошибка не отменяет приём).
|
||||
// tenantScope: Scope тенанта (TenantDbContext построен на схеме тенанта).
|
||||
// tenant: Тенант канала (для лога аудита).
|
||||
// request: Сообщение PushMessage.
|
||||
// ct: Токен отмены.
|
||||
private async Task SavePreviewSafelyAsync(
|
||||
AsyncServiceScope tenantScope,
|
||||
TenantRecordDto tenant,
|
||||
PushMessageRequest request,
|
||||
CancellationToken ct)
|
||||
{
|
||||
try
|
||||
{
|
||||
DialogsService dialogs = tenantScope.ServiceProvider.GetRequiredService<DialogsService>();
|
||||
DateTimeOffset? msgAt = request.HasMsgAt ? DateTimeOffset.FromUnixTimeMilliseconds(request.MsgAt) : null;
|
||||
await dialogs.SavePreviewAsync(
|
||||
request.DialogId,
|
||||
request.HasMsgId ? request.MsgId : null,
|
||||
request.Text,
|
||||
msgAt,
|
||||
ct).ConfigureAwait(false);
|
||||
}
|
||||
catch (OperationCanceledException)
|
||||
{
|
||||
throw;
|
||||
}
|
||||
catch (Exception exception)
|
||||
{
|
||||
// Превью — вторичная запись: приём не затронут (лог дебага, не ошибка RPC).
|
||||
logger.LogDebug(exception, "PushMessage {TenantId}: превью не сохранено (приём не затронут)", tenant.Id);
|
||||
}
|
||||
}
|
||||
|
||||
// Публикует SSE system_status (каждый репорт) и тосты на переходах connected.
|
||||
// tenantId: Тенант канала (реестровый Guid).
|
||||
// previous: Предыдущий снимок из KV (null — первый репорт).
|
||||
// current: Текущий снимок репорта.
|
||||
private void PublishStatusEvents(
|
||||
Guid tenantId,
|
||||
TgReportedStatus? previous,
|
||||
TgReportedStatus current)
|
||||
{
|
||||
broker.Publish(tenantId, SystemStatusEventType, current);
|
||||
if (previous is null)
|
||||
{
|
||||
// Первый репорт после старта сервиса: переходов нет, статус фронт получит по system_status.
|
||||
return;
|
||||
}
|
||||
|
||||
if (!previous.Connected && current.Connected)
|
||||
{
|
||||
PublishToast(tenantId, ConnectedToastText, ConnectedToastIcon);
|
||||
}
|
||||
else if (previous.Connected && !current.Connected)
|
||||
{
|
||||
PublishToast(tenantId, DisconnectedToastText, DisconnectedToastIcon);
|
||||
}
|
||||
}
|
||||
|
||||
// Публикует SSE-тост в канал тенанта (без подписчиков — no-op, Ruling 5).
|
||||
// tenantId: Тенант-получатель.
|
||||
// text: Текст тоста.
|
||||
// icon: Иконка тоста (набор Icon.vue фронта).
|
||||
private void PublishToast(
|
||||
Guid tenantId,
|
||||
string text,
|
||||
string icon)
|
||||
{
|
||||
broker.Publish(tenantId, ToastEventType, new { text, icon });
|
||||
}
|
||||
|
||||
// Снимок предыдущего статуса из KV tgStatus (нет записи/битый JSON — null).
|
||||
// settings: KV-хранилище настроек схемы тенанта.
|
||||
// ct: Токен отмены.
|
||||
// Возвращает: Предыдущий снимок либо null.
|
||||
private async Task<TgReportedStatus?> ReadPreviousStatusAsync(ISettingsStore settings, CancellationToken ct)
|
||||
{
|
||||
SettingValue? stored = await settings.GetAsync(SettingsKeys.TgStatus, ct).ConfigureAwait(false);
|
||||
if (stored is null)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
return JsonSerializer.Deserialize<TgReportedStatus>(stored.ValueJson, StatusJsonOptions);
|
||||
}
|
||||
catch (JsonException exception)
|
||||
{
|
||||
logger.LogWarning(exception, "Аудит: KV tgStatus повреждён — переходы фаз не определяются");
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
// Маппит запрос ReportStatus в снимок KV (account живёт отдельным ключом tgAccount).
|
||||
// request: Запрос ReportStatus.
|
||||
// Возвращает: Снимок статуса.
|
||||
private static TgReportedStatus ToReportedStatus(ReportStatusRequest request) => new()
|
||||
{
|
||||
Phase = request.Phase,
|
||||
Connected = request.Connected,
|
||||
Listener = request.Listener,
|
||||
Error = request.HasError ? request.Error : null,
|
||||
QrUrl = request.HasQrUrl ? request.QrUrl : null,
|
||||
};
|
||||
|
||||
// Сериализует снимок в JSON (camelCase, конвенция value_json).
|
||||
// status: Снимок статуса.
|
||||
// Возвращает: JSON-строка.
|
||||
private static string ToJson(TgReportedStatus status) => JsonSerializer.Serialize(status, StatusJsonOptions);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user