Перевести входящий поток источников на generic-контракт

Добавлен sources.proto с PushSource; приём в ядре вынесен в SourceIngressGrpcService с SourceProtoMapper и ISourceIngestObserver, тенант определяется IngressTenantResolver. Из telegram.proto удалён PushMessage, telegram-сервис шлёт generic-записи, превью сохраняет TelegramSourceIngestObserver.
This commit is contained in:
Rustam Khalimov
2026-09-11 16:32:49 +03:00
parent 3327bf48b0
commit 6d074834a7
39 changed files with 1312 additions and 532 deletions
@@ -2,6 +2,7 @@ using Deal.Grpc.Hosting.Interceptors;
using Deal.Grpc.Hosting.Models;
using Deal.Grpc.Hosting.Options;
using Deal.Grpc.Hosting.Services;
using Deal.Grpc.Sources;
using Deal.Grpc.Telegram;
using Deal.Telegram.Sessions;
using Grpc.Core;
@@ -22,7 +23,8 @@ public sealed class CoreIngressClient : ICoreIngressClient
private readonly ILogger<CoreIngressClient> _logger;
private readonly MtlsCertificates? _mtlsCertificates;
private readonly object _channelGate = new();
private IngressService.IngressServiceClient? _client;
private IngressService.IngressServiceClient? _ingressClient;
private SourceIngressService.SourceIngressServiceClient? _sourceClient;
private GrpcChannel? _channel;
/// <summary>
@@ -42,19 +44,19 @@ public sealed class CoreIngressClient : ICoreIngressClient
}
/// <inheritdoc />
public async Task<PushMessageReply> PushMessageAsync(
public async Task<PushSourceReply> PushSourceAsync(
string tenantId,
PushMessageRequest message,
PushSourceRequest request,
CancellationToken cancellationToken)
{
IngressService.IngressServiceClient client = GetClient();
SourceIngressService.SourceIngressServiceClient client = GetSourceClient();
try
{
return await client.PushMessageAsync(message, CallOptions(tenantId)).ResponseAsync.WaitAsync(cancellationToken).ConfigureAwait(false);
return await client.PushSourceAsync(request, CallOptions(tenantId)).ResponseAsync.WaitAsync(cancellationToken).ConfigureAwait(false);
}
catch (Exception exception) when (exception is not OperationCanceledException)
{
throw Fail("PushMessage", exception);
throw Fail("PushSource", exception);
}
}
@@ -64,7 +66,7 @@ public sealed class CoreIngressClient : ICoreIngressClient
IReadOnlyList<DialogEntry> entries,
CancellationToken cancellationToken)
{
IngressService.IngressServiceClient client = GetClient();
IngressService.IngressServiceClient client = GetIngressClient();
var request = new SyncDialogsRequest();
request.Entries.AddRange(entries);
try
@@ -78,31 +80,44 @@ public sealed class CoreIngressClient : ICoreIngressClient
}
}
// Лениво создаёт gRPC-канал к ядру (адрес неизменен на время жизни процесса).
private IngressService.IngressServiceClient GetClient()
// Клиент SourceIngress (PushSource); канал к ядру создаётся лениво.
private SourceIngressService.SourceIngressServiceClient GetSourceClient()
{
lock (_channelGate)
{
if (_client is null)
{
if (_mtlsCertificates is not null)
{
_channel = GrpcChannel.ForAddress(
_options.IngressEndpoint,
new GrpcChannelOptions { HttpHandler = _mtlsCertificates.CreateClientHttpHandler() });
}
else
{
_channel = GrpcChannel.ForAddress(_options.IngressEndpoint);
}
_client = new IngressService.IngressServiceClient(_channel);
}
return _client;
EnsureChannelLocked();
return _sourceClient!;
}
}
// Клиент Ingress (SyncDialogs); канал к ядру создаётся лениво.
private IngressService.IngressServiceClient GetIngressClient()
{
lock (_channelGate)
{
EnsureChannelLocked();
return _ingressClient!;
}
}
// Лениво создаёт gRPC-канал к ядру (адрес неизменен на время жизни процесса).
private void EnsureChannelLocked()
{
if (_channel is not null)
{
return;
}
_channel = _mtlsCertificates is not null
? GrpcChannel.ForAddress(
_options.IngressEndpoint,
new GrpcChannelOptions { HttpHandler = _mtlsCertificates.CreateClientHttpHandler() })
: GrpcChannel.ForAddress(_options.IngressEndpoint);
_ingressClient = new IngressService.IngressServiceClient(_channel);
_sourceClient = new SourceIngressService.SourceIngressServiceClient(_channel);
}
private CallOptions CallOptions(string tenantId)
{
var metadata = new Metadata
@@ -1,3 +1,4 @@
using Deal.Grpc.Sources;
using Deal.Grpc.Telegram;
namespace Deal.Telegram.Core;
@@ -8,15 +9,15 @@ namespace Deal.Telegram.Core;
public interface ICoreIngressClient
{
/// <summary>
/// Отправляет сообщение диалога в ядро.
/// Отправляет запись источника в ядро.
/// </summary>
/// <param name="tenantId">Id тенанта.</param>
/// <param name="message">Сообщение в контракте PushMessageRequest.</param>
/// <param name="request">Запись в контракте PushSourceRequest.</param>
/// <param name="cancellationToken">Отмена вызова.</param>
/// <returns>Ответ ядра (accepted/duplicate — дубль dialog+msgId в очереди не растёт).</returns>
public Task<PushMessageReply> PushMessageAsync(
/// <returns>Ответ ядра (accepted/duplicate — дубль источника в очереди не растёт).</returns>
public Task<PushSourceReply> PushSourceAsync(
string tenantId,
PushMessageRequest message,
PushSourceRequest request,
CancellationToken cancellationToken);
/// <summary>
@@ -1,5 +1,5 @@
using System.Collections.Concurrent;
using Deal.Grpc.Telegram;
using Deal.Grpc.Sources;
using Deal.Telegram.Core;
using Deal.Telegram.Sessions;
using Deal.Telegram.Telegram;
@@ -54,7 +54,7 @@ public sealed class BackfillService
/// Создаёт службу backfill'а.
/// </summary>
/// <param name="sessionFarm">Пул сессий тенантов (read-операции диалогов).</param>
/// <param name="ingress">Канал в ядро (PushMessage).</param>
/// <param name="ingress">Канал в ядро (PushSource).</param>
/// <param name="pacer">Анти-бан-паузы (реальный — случайные, тесты — фейк).</param>
/// <param name="logger">Логгер.</param>
public BackfillService(
@@ -168,8 +168,8 @@ public sealed class BackfillService
int processed = 0;
foreach (TelegramMessage message in messages.Reverse())
{
PushMessageReply reply = await _ingress
.PushMessageAsync(tenantId, DialogProtoMapper.ToPushRequest(message), cancellationToken)
PushSourceReply reply = await _ingress
.PushSourceAsync(tenantId, DialogProtoMapper.ToSourceRequest(message), cancellationToken)
.ConfigureAwait(false);
if (reply.Accepted && !reply.Duplicate)
{
@@ -1,4 +1,5 @@
using System.Globalization;
using Deal.Grpc.Sources;
using Deal.Grpc.Telegram;
using Deal.Telegram.Telegram;
@@ -9,6 +10,10 @@ namespace Deal.Telegram.Dialogs;
/// </summary>
public static class DialogProtoMapper
{
private const string TelegramKind = "telegram";
private const string HueKey = "hue";
/// <summary>
/// Нейтральный диалог → DialogEntry контракта
/// </summary>
@@ -38,19 +43,27 @@ public static class DialogProtoMapper
};
/// <summary>
/// Нейтральное сообщение → PushMessageRequest ингресса.
/// Нейтральное сообщение → PushSourceRequest ингресса.
/// </summary>
/// <param name="message">Сообщение диалога (непустой текст).</param>
/// <returns>Запрос Ingress.PushMessage с канальными полями и дубль-гвардом msg_id.</returns>
public static PushMessageRequest ToPushRequest(TelegramMessage message)
/// <returns>Запрос SourceIngress.PushSource с источником telegram и содержимым.</returns>
public static PushSourceRequest ToSourceRequest(TelegramMessage message)
=> new()
{
DialogId = message.DialogId,
ChannelName = message.DialogName,
ChannelHandle = message.DialogHandle,
ChannelHue = DialogHue.Compute(message.DialogId, message.DialogName),
MsgId = message.Id,
Text = message.Text,
MsgAt = message.DateMs,
Source = new SourceRefProto
{
Kind = TelegramKind,
ExternalId = message.Id.ToString(CultureInfo.InvariantCulture),
OriginRef = message.DialogId,
DisplayName = message.DialogName,
Author = message.DialogName,
ReceivedAt = message.DateMs,
Extra = { [HueKey] = DialogHue.Compute(message.DialogId, message.DialogName) },
},
Content = new SourceContentProto
{
Text = message.Text,
Author = message.DialogName,
},
};
}
@@ -1,4 +1,4 @@
using Deal.Grpc.Telegram;
using Deal.Grpc.Sources;
using Deal.Telegram.Core;
using Deal.Telegram.Sessions;
using Deal.Telegram.Telegram;
@@ -20,7 +20,7 @@ public sealed class RealtimeListener
/// </summary>
/// <param name="session">Ready-сессия тенанта (события сообщений её клиента).</param>
/// <param name="catalog">Зеркало мониторинга.</param>
/// <param name="ingress">Канал в ядро (PushMessage).</param>
/// <param name="ingress">Канал в ядро (PushSource).</param>
/// <param name="logger">Логгер.</param>
public RealtimeListener(
TenantSession session,
@@ -52,7 +52,7 @@ public sealed class RealtimeListener
_session.SetListenerActive(false);
}
// Обработчик входящего сообщения: фильтр по зеркалу → PushMessage в ядро → mark-as-read.
// Обработчик входящего сообщения: фильтр по зеркалу → PushSource в ядро → mark-as-read.
// Сбой отправки не роняет realtime: сообщение остаётся непрочитанным и догоняется realtime_sweep.
// message: Входящее сообщение аккаунта.
private async Task OnMessageReceivedAsync(TelegramMessage message)
@@ -64,8 +64,8 @@ public sealed class RealtimeListener
return;
}
PushMessageReply reply = await _ingress
.PushMessageAsync(_session.TenantId, DialogProtoMapper.ToPushRequest(message), CancellationToken.None)
PushSourceReply reply = await _ingress
.PushSourceAsync(_session.TenantId, DialogProtoMapper.ToSourceRequest(message), CancellationToken.None)
.ConfigureAwait(false);
if (reply.Accepted)
{
@@ -1,3 +1,4 @@
using Deal.Grpc.Sources;
using Deal.Grpc.Telegram;
using Deal.Telegram.Core;
using Deal.Telegram.Sessions;
@@ -40,7 +41,7 @@ public sealed class RealtimeSweep
/// </summary>
/// <param name="sessionFarm">Пул сессий тенантов.</param>
/// <param name="catalog">Зеркало каталога/мониторинга (актуализируется ответом SyncDialogs).</param>
/// <param name="ingress">Канал в ядро (PushMessage/SyncDialogs).</param>
/// <param name="ingress">Канал в ядро (PushSource/SyncDialogs).</param>
/// <param name="logger">Логгер.</param>
public RealtimeSweep(
SessionFarm sessionFarm,
@@ -145,8 +146,8 @@ public sealed class RealtimeSweep
int added = 0;
foreach (TelegramMessage message in messages.Reverse())
{
PushMessageReply reply = await _ingress
.PushMessageAsync(tenantId, DialogProtoMapper.ToPushRequest(message), cancellationToken)
PushSourceReply reply = await _ingress
.PushSourceAsync(tenantId, DialogProtoMapper.ToSourceRequest(message), cancellationToken)
.ConfigureAwait(false);
if (reply.Accepted && !reply.Duplicate)
{
@@ -27,7 +27,7 @@ public sealed class RealtimeMonitorService : BackgroundService
/// </summary>
/// <param name="sessionFarm">Пул сессий тенантов.</param>
/// <param name="catalog">Зеркало мониторинга (фильтр сообщений).</param>
/// <param name="ingress">Канал в ядро (PushMessage).</param>
/// <param name="ingress">Канал в ядро (PushSource).</param>
/// <param name="loggerFactory">Фабрика логгеров (логгеры listener'ов).</param>
public RealtimeMonitorService(
SessionFarm sessionFarm,
@@ -53,7 +53,7 @@ public static class SessionErrorMessages
public const string TelegramUnavailable = "Telegram недоступен — повторите попытку позже";
/// <summary>
/// Ядро (gRPC-ингресс) недоступно — PushMessage/SyncDialogs не доставлены (UNAVAILABLE).
/// Ядро (gRPC-ингресс) недоступно — PushSource/SyncDialogs не доставлены (UNAVAILABLE).
/// </summary>
public const string IngressUnavailable = "Ядро недоступно — повторите попытку позже";
@@ -92,7 +92,7 @@ public interface ISessionClient : IAsyncDisposable
/// <param name="dialogId">Подписанный id диалога (каналы "-100…", группы "-…", личные "+…").</param>
/// <param name="limit">Сколько последних сообщений запросить.</param>
/// <param name="cancellationToken">Отмена операции.</param>
/// <returns>Сообщения диалога (с канальными полями для PushMessage).</returns>
/// <returns>Сообщения диалога (с канальными полями для PushSource).</returns>
public Task<IReadOnlyList<TelegramMessage>> GetMessagesAsync(
string dialogId,
int limit,
@@ -12,8 +12,8 @@ public sealed record TelegramMessage
/// <param name="id">Id сообщения в Telegram (дубль-гвард dialog+msgId ядра).</param>
/// <param name="text">Текст сообщения (непустой).</param>
/// <param name="dateMs">Время сообщения, epoch-ms.</param>
/// <param name="dialogName">Имя диалога (title/first_name) для PushMessage.channel_name.</param>
/// <param name="dialogHandle">Username диалога для PushMessage.channel_handle.</param>
/// <param name="dialogName">Имя диалога (title/first_name) для source.display_name.</param>
/// <param name="dialogHandle">Username диалога.</param>
public TelegramMessage(
string dialogId,
int id,
@@ -51,7 +51,7 @@ public sealed record TelegramMessage
public long DateMs { get; }
/// <summary>
/// Имя диалога (title/first_name) — для PushMessage.channel_name.
/// Имя диалога (title/first_name) — для source.display_name.
/// </summary>
public string DialogName { get; }