This commit is contained in:
kislovdm
2026-04-02 21:36:56 +03:00
parent c79a4d4f6d
commit 4c08ebd68b
2 changed files with 264 additions and 33 deletions
+236 -33
View File
@@ -1,3 +1,4 @@
using System.Collections.Concurrent;
using Telegram.Bot; using Telegram.Bot;
using Telegram.Bot.Types; using Telegram.Bot.Types;
using Telegram.Bot.Types.Enums; using Telegram.Bot.Types.Enums;
@@ -8,11 +9,17 @@ namespace Telegram2VkBot;
public sealed class ForwardWorker : BackgroundService public sealed class ForwardWorker : BackgroundService
{ {
private static readonly TimeSpan AlbumDebounce = TimeSpan.FromMilliseconds(950);
private readonly TelegramOptions _telegram; private readonly TelegramOptions _telegram;
private readonly VkOptions _vk; private readonly VkOptions _vk;
private readonly VkApiClient _vkApi; private readonly VkApiClient _vkApi;
private readonly MappingRepository _repo; private readonly MappingRepository _repo;
private readonly ILogger<ForwardWorker> _logger; private readonly ILogger<ForwardWorker> _logger;
private readonly SemaphoreSlim _updateGate = new(1, 1);
private readonly ConcurrentDictionary<(long ChatId, string MediaGroupId), AlbumBuffer> _albumBuffers = new();
private CancellationToken _appStopping;
public ForwardWorker( public ForwardWorker(
IOptions<TelegramOptions> telegram, IOptions<TelegramOptions> telegram,
@@ -30,6 +37,8 @@ public sealed class ForwardWorker : BackgroundService
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{ {
_appStopping = stoppingToken;
if (string.IsNullOrWhiteSpace(_telegram.BotToken)) if (string.IsNullOrWhiteSpace(_telegram.BotToken))
throw new ArgumentException("TELEGRAM__BotToken (или TELEGRAM:BotToken) пуст — проверьте .env / переменные окружения."); throw new ArgumentException("TELEGRAM__BotToken (или TELEGRAM:BotToken) пуст — проверьте .env / переменные окружения.");
if (_telegram.ChannelId == 0) if (_telegram.ChannelId == 0)
@@ -54,29 +63,39 @@ public sealed class ForwardWorker : BackgroundService
_logger.LogInformation("Start Telegram polling for channel {ChannelId}", _telegram.ChannelId); _logger.LogInformation("Start Telegram polling for channel {ChannelId}", _telegram.ChannelId);
using var semaphore = new SemaphoreSlim(1, 1);
bot.StartReceiving( bot.StartReceiving(
async (client, update, ct) => async (client, update, ct) =>
{ {
await semaphore.WaitAsync(ct); await _updateGate.WaitAsync(ct);
try try
{ {
await HandleUpdateAsync(update, bot, ct); await HandleUpdateAsync(update, bot, ct);
} }
finally finally
{ {
semaphore.Release(); _updateGate.Release();
} }
}, },
HandleErrorAsync, HandleErrorAsync,
receiverOptions, receiverOptions,
stoppingToken); stoppingToken);
// StartReceiving не блокирует, поэтому держим сервис живым.
await Task.Delay(Timeout.InfiniteTimeSpan, stoppingToken); await Task.Delay(Timeout.InfiniteTimeSpan, stoppingToken);
} }
public override void Dispose()
{
_updateGate.Dispose();
foreach (var kv in _albumBuffers)
{
kv.Value.DebounceCts?.Cancel();
kv.Value.DebounceCts?.Dispose();
}
_albumBuffers.Clear();
base.Dispose();
}
private async Task HandleUpdateAsync(Update update, TelegramBotClient bot, CancellationToken ct) private async Task HandleUpdateAsync(Update update, TelegramBotClient bot, CancellationToken ct)
{ {
if (update.EditedChannelPost is { } edited && edited.Chat.Id == _telegram.ChannelId) if (update.EditedChannelPost is { } edited && edited.Chat.Id == _telegram.ChannelId)
@@ -93,14 +112,102 @@ public sealed class ForwardWorker : BackgroundService
} }
private async Task HandleChannelPostAsync(TelegramBotClient bot, Message message, CancellationToken ct) private async Task HandleChannelPostAsync(TelegramBotClient bot, Message message, CancellationToken ct)
{
if (message.MediaGroupId is { } groupId)
{
EnqueueAlbumPartAndScheduleFlush(bot, message, groupId);
return;
}
await PostSingleChannelMessageAsync(bot, message, ct);
}
private void EnqueueAlbumPartAndScheduleFlush(TelegramBotClient bot, Message message, string mediaGroupId)
{
var key = (message.Chat.Id, mediaGroupId);
var buffer = _albumBuffers.GetOrAdd(key, _ => new AlbumBuffer());
lock (buffer.Sync)
{
if (buffer.Messages.TrueForAll(m => m.MessageId != message.MessageId))
buffer.Messages.Add(message);
buffer.DebounceCts?.Cancel();
buffer.DebounceCts?.Dispose();
buffer.DebounceCts = CancellationTokenSource.CreateLinkedTokenSource(_appStopping);
var debounceToken = buffer.DebounceCts.Token;
_logger.LogDebug(
"Альбом media_group_id={MediaGroupId}: частей в буфере={Count}",
mediaGroupId,
buffer.Messages.Count);
_ = FlushAlbumAfterDebounceAsync(bot, key, debounceToken);
}
}
private async Task FlushAlbumAfterDebounceAsync(
TelegramBotClient bot,
(long ChatId, string MediaGroupId) key,
CancellationToken debounceToken)
{
try
{
await Task.Delay(AlbumDebounce, debounceToken).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
return;
}
await _updateGate.WaitAsync(debounceToken).ConfigureAwait(false);
try
{
await FlushAlbumBufferCoreAsync(bot, key, debounceToken).ConfigureAwait(false);
}
finally
{
_updateGate.Release();
}
}
private async Task FlushAlbumBufferCoreAsync(
TelegramBotClient bot,
(long ChatId, string MediaGroupId) key,
CancellationToken ct)
{
if (!_albumBuffers.TryRemove(key, out var buffer))
return;
buffer.DebounceCts?.Dispose();
buffer.DebounceCts = null;
List<Message> batch;
lock (buffer.Sync)
{
batch = buffer.Messages.OrderBy(m => m.MessageId).ToList();
}
if (batch.Count == 0)
return;
_logger.LogInformation(
"Альбом media_group_id={MediaGroupId}: публикация одного поста ВК по {Count} частям Telegram (message ids: {Ids})",
key.MediaGroupId,
batch.Count,
string.Join(",", batch.Select(m => m.MessageId)));
await PostChannelMessagesAsSingleVkPostAsync(bot, batch, ct);
}
private async Task PostSingleChannelMessageAsync(TelegramBotClient bot, Message message, CancellationToken ct)
{ {
var ownerId = -Math.Abs(_vk.GroupId); var ownerId = -Math.Abs(_vk.GroupId);
var telegramMessageId = message.MessageId; var telegramMessageId = message.MessageId;
var text = ExtractTelegramText(message); var text = ExtractTelegramText(message);
var attachments = await ExtractPhotoAttachmentIfAnyAsync(bot, message, ct); var attachments = await ExtractPhotoAttachmentsFromMessagesAsync(bot, new[] { message }, ct);
// Если это не текст и не фото, просто пропускаем (сложно маппить видео/документы в VK без доп. логики).
if (string.IsNullOrWhiteSpace(text) && string.IsNullOrWhiteSpace(attachments)) if (string.IsNullOrWhiteSpace(text) && string.IsNullOrWhiteSpace(attachments))
{ {
_logger.LogInformation("Skip message {MessageId}: no text/photo", telegramMessageId); _logger.LogInformation("Skip message {MessageId}: no text/photo", telegramMessageId);
@@ -117,17 +224,58 @@ public sealed class ForwardWorker : BackgroundService
var postId = await _vkApi.WallPostAsync(ownerId, text ?? string.Empty, attachments, ct); var postId = await _vkApi.WallPostAsync(ownerId, text ?? string.Empty, attachments, ct);
await _repo.UpsertAsync( await _repo.UpsertAsync(
telegramChatId: _telegram.ChannelId, _telegram.ChannelId,
telegramMessageId: telegramMessageId, telegramMessageId,
vkOwnerId: ownerId, ownerId,
vkPostId: postId, postId,
vkMessage: text, text,
vkAttachments: attachments, attachments,
ct); ct);
_logger.LogInformation("Posted new VK post {PostId} for Telegram {MessageId}", postId, telegramMessageId); _logger.LogInformation("Posted new VK post {PostId} for Telegram {MessageId}", postId, telegramMessageId);
} }
private async Task PostChannelMessagesAsSingleVkPostAsync(TelegramBotClient bot, IReadOnlyList<Message> batch, CancellationToken ct)
{
var ownerId = -Math.Abs(_vk.GroupId);
var text = ExtractTelegramTextFromAlbum(batch);
var attachments = await ExtractPhotoAttachmentsFromMessagesAsync(bot, batch, ct);
if (string.IsNullOrWhiteSpace(text) && string.IsNullOrWhiteSpace(attachments))
{
_logger.LogInformation("Skip album: no text/photo after merge (ids {Ids})", string.Join(",", batch.Select(m => m.MessageId)));
return;
}
_logger.LogInformation(
"VK wall.post (альбом): owner_id={OwnerId}, частей={Count}, textLen={TextLen}, textPreview={TextPreview}, attachments={Attachments}",
ownerId,
batch.Count,
(text ?? string.Empty).Length,
TruncateForLog(text, 200),
attachments ?? "(none)");
var postId = await _vkApi.WallPostAsync(ownerId, text ?? string.Empty, attachments, ct);
foreach (var msg in batch)
{
await _repo.UpsertAsync(
_telegram.ChannelId,
msg.MessageId,
ownerId,
postId,
text,
attachments,
ct);
}
_logger.LogInformation(
"Posted single VK post {PostId} for Telegram album media_group_id={Group} (message ids {Ids})",
postId,
batch[0].MediaGroupId,
string.Join(",", batch.Select(m => m.MessageId)));
}
private async Task HandleChannelPostEditAsync(TelegramBotClient bot, Message editedMessage, CancellationToken ct) private async Task HandleChannelPostEditAsync(TelegramBotClient bot, Message editedMessage, CancellationToken ct)
{ {
var telegramMessageId = editedMessage.MessageId; var telegramMessageId = editedMessage.MessageId;
@@ -137,8 +285,15 @@ public sealed class ForwardWorker : BackgroundService
var extractedText = ExtractTelegramText(editedMessage); var extractedText = ExtractTelegramText(editedMessage);
// Если фото в апдейте отсутствует, то оставляем старые вложения, чтобы не "стереть" медиа на стороне VK. var extractedAttachments = await ExtractPhotoAttachmentsFromMessagesAsync(bot, new[] { editedMessage }, ct);
string? extractedAttachments = await ExtractPhotoAttachmentIfAnyAsync(bot, editedMessage, ct);
// У апдейта альбома в Telegram приходит одно сообщение с одной фото; не затираем все вложения ВК одним файлом.
if (editedMessage.MediaGroupId is not null && existing is { VkAttachments: { } prevAtt })
{
if (CountAttachmentTokens(prevAtt) > CountAttachmentTokens(extractedAttachments))
extractedAttachments = prevAtt;
}
var attachmentsToSend = !string.IsNullOrWhiteSpace(extractedAttachments) var attachmentsToSend = !string.IsNullOrWhiteSpace(extractedAttachments)
? extractedAttachments ? extractedAttachments
: existing?.VkAttachments; : existing?.VkAttachments;
@@ -181,11 +336,69 @@ public sealed class ForwardWorker : BackgroundService
attachmentsToSend ?? "(none)"); attachmentsToSend ?? "(none)");
await _vkApi.WallEditAsync(existing.Value.VkOwnerId, existing.Value.VkPostId, messageToSendFinal, attachmentsToSend, ct); await _vkApi.WallEditAsync(existing.Value.VkOwnerId, existing.Value.VkPostId, messageToSendFinal, attachmentsToSend, ct);
await _repo.UpsertAsync(_telegram.ChannelId, telegramMessageId, existing.Value.VkOwnerId, existing.Value.VkPostId, messageToSendFinal, attachmentsToSend, ct);
var idsToTouch = await _repo.GetTelegramMessageIdsForVkPostAsync(_telegram.ChannelId, existing.Value.VkPostId, ct);
if (idsToTouch.Count == 0)
idsToTouch = new[] { telegramMessageId };
foreach (var mid in idsToTouch)
{
await _repo.UpsertAsync(
_telegram.ChannelId,
mid,
existing.Value.VkOwnerId,
existing.Value.VkPostId,
messageToSendFinal,
attachmentsToSend,
ct);
}
_logger.LogInformation("Edited VK post {PostId} for Telegram {MessageId}", existing.Value.VkPostId, telegramMessageId); _logger.LogInformation("Edited VK post {PostId} for Telegram {MessageId}", existing.Value.VkPostId, telegramMessageId);
} }
private static int CountAttachmentTokens(string? vkAttachmentsCsv)
{
if (string.IsNullOrWhiteSpace(vkAttachmentsCsv))
return 0;
return vkAttachmentsCsv.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries).Length;
}
private static string? ExtractTelegramTextFromAlbum(IReadOnlyList<Message> messages)
{
foreach (var message in messages.OrderBy(m => m.MessageId))
{
var t = ExtractTelegramText(message);
if (!string.IsNullOrWhiteSpace(t))
return t;
}
return null;
}
private async Task<string?> ExtractPhotoAttachmentsFromMessagesAsync(
TelegramBotClient bot,
IReadOnlyList<Message> messages,
CancellationToken ct)
{
var ordered = messages.OrderBy(m => m.MessageId).ToList();
var parts = new List<string>();
foreach (var message in ordered)
{
if (message.Photo is not { Length: > 0 })
continue;
var best = message.Photo[^1];
await using var ms = new MemoryStream();
await bot.GetInfoAndDownloadFile(best.FileId, ms, ct).ConfigureAwait(false);
var att = await _vkApi.UploadPhotoAndGetAttachmentAsync(ms.ToArray(), ct).ConfigureAwait(false);
if (!string.IsNullOrWhiteSpace(att))
parts.Add(att);
}
return parts.Count == 0 ? null : string.Join(",", parts);
}
private static string? ExtractTelegramText(Message message) private static string? ExtractTelegramText(Message message)
{ {
if (!string.IsNullOrWhiteSpace(message.Text)) if (!string.IsNullOrWhiteSpace(message.Text))
@@ -195,28 +408,12 @@ public sealed class ForwardWorker : BackgroundService
return null; return null;
} }
private async Task<string?> ExtractPhotoAttachmentIfAnyAsync(TelegramBotClient bot, Message message, CancellationToken ct)
{
if (message.Photo == null || message.Photo.Length == 0)
return null;
// Telegram предоставляет несколько размеров одного фото. Берем самый большой.
var best = message.Photo[^1];
await using var ms = new MemoryStream();
// В Telegram.Bot v22 методы называются без Async-суффикса.
await bot.GetInfoAndDownloadFile(best.FileId, ms, ct);
return await _vkApi.UploadPhotoAndGetAttachmentAsync(ms.ToArray(), ct);
}
private static Task HandleErrorAsync( private static Task HandleErrorAsync(
ITelegramBotClient client, ITelegramBotClient client,
Exception exception, Exception exception,
HandleErrorSource errorSource, HandleErrorSource errorSource,
CancellationToken cancellationToken) CancellationToken cancellationToken)
{ {
// Логирование делаем в Worker через исключения, но чтобы не терять ошибки polling'а — просто возвращаем CompletedTask.
return Task.CompletedTask; return Task.CompletedTask;
} }
@@ -228,5 +425,11 @@ public sealed class ForwardWorker : BackgroundService
return text; return text;
return string.Concat(text.AsSpan(0, maxLen), "…(len=", text.Length.ToString(), ")"); return string.Concat(text.AsSpan(0, maxLen), "…(len=", text.Length.ToString(), ")");
} }
}
private sealed class AlbumBuffer
{
public object Sync { get; } = new();
public List<Message> Messages { get; } = new();
public CancellationTokenSource? DebounceCts;
}
}
+28
View File
@@ -75,6 +75,34 @@ LIMIT 1;";
return null; return null;
} }
/// <summary>Все telegram message_id, привязанные к одному посту ВК (альбом).</summary>
public async Task<IReadOnlyList<int>> GetTelegramMessageIdsForVkPostAsync(
long telegramChatId,
int vkPostId,
CancellationToken ct)
{
await using var con = new SqliteConnection(ConnectionString);
await con.OpenAsync(ct);
var sql = @"
SELECT telegram_message_id
FROM vk_mappings
WHERE telegram_chat_id = $telegram_chat_id AND vk_post_id = $vk_post_id
ORDER BY telegram_message_id;";
await using var cmd = con.CreateCommand();
cmd.CommandText = sql;
cmd.Parameters.AddWithValue("$telegram_chat_id", telegramChatId);
cmd.Parameters.AddWithValue("$vk_post_id", vkPostId);
var ids = new List<int>();
await using var reader = await cmd.ExecuteReaderAsync(ct);
while (await reader.ReadAsync(ct))
ids.Add(reader.GetInt32(0));
return ids;
}
public async Task UpsertAsync( public async Task UpsertAsync(
long telegramChatId, long telegramChatId,
int telegramMessageId, int telegramMessageId,