WebRtc модуль и документация, импорт телеграма
This commit is contained in:
@@ -0,0 +1,52 @@
|
|||||||
|
# Модуль Импорта Telegram (Telegram Import Module) — Knot Messager
|
||||||
|
|
||||||
|
Модуль предназначен для бесшовного переноса истории переписки из Telegram (HTML Export) в защищенную среду Knot Messager. Основной упор сделан на сохранение контекста, вложений и связей между пользователями через систему мэппинга.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🏗 Архитектура Процесса
|
||||||
|
|
||||||
|
Импорт реализован как двухфазная операция для обеспечения максимальной точности данных и стабильности сервера:
|
||||||
|
|
||||||
|
### Фаза 1: Предварительный Анализ (Pre-Analysis)
|
||||||
|
На этом этапе система:
|
||||||
|
1. **Extracts Names**: Сканирует ZIP-архив и вытаскивает список всех отправителей (`from_name`).
|
||||||
|
2. **Validates Consistency**: Проверяет целостность HTML-файлов и наличие папок с медиа.
|
||||||
|
3. **Analyzes Policies**: Сравнивает контент из Telegram с [SystemSettings](file:///e:/GIT/forkmessager/backend/src/Shared/Knot.Shared.Kernel/Configuration/ISettingsService.cs#13-17) вашего узла (разрешены ли медиа, голосовые, опросы).
|
||||||
|
4. **Returns Conflicts**: Выдает фронтенду список флагов (например, `Media.Blocked`), чтобы предупредить пользователя о потере данных из-за политик сервера.
|
||||||
|
|
||||||
|
### Фаза 2: Фоновое Исполнение (Background Execution)
|
||||||
|
После получения подтверждения и мэппинга имен:
|
||||||
|
1. **Job Queuing**: Создается задача импорта, возвращается `JobId`.
|
||||||
|
2. **Worker Processing**: Система в фоновом потоке ([TelegramImportWorker](file:///e:/GIT/forkmessager/backend/src/Modules/TelegramImport/Infrastructure/Background/TelegramImportWorker.cs#21-98)) начинает парсинг и сохранение сообщений порциями (chunking), чтобы не перегружать RAM.
|
||||||
|
3. **Real-time Progress**: Статус задачи обновляется в [ImportJobStore](file:///e:/GIT/forkmessager/backend/src/Modules/TelegramImport/Infrastructure/Background/ImportJobStore.cs#31-38), что позволяет фронтенду отображать прогресс-бар (0-100%).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 👥 Мэппинг Пользователей (User Mapping)
|
||||||
|
|
||||||
|
Одной из ключевых функций модуля является возможность передать `Dictionary<string, Guid> Ranking`.
|
||||||
|
* **Имя из Telegram**: Строка (например, "Ivan_HR").
|
||||||
|
* **Knot UserId**: UUID существующего контакта в модуле [Relations](file:///e:/GIT/forkmessager/backend/src/Modules/Relations/Infrastructure/Persistence/RelationsDbContext.cs#17-22).
|
||||||
|
* **Результат**: При импорте все сообщения от "Ivan_HR" будут автоматически привязаны к вашему реальному контакту с его аватаром и настройками.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🔌 API Эндпоинты
|
||||||
|
|
||||||
|
* `POST /api/import/analyze` — Загрузка ZIP, возврат списка имен и конфликтов политик.
|
||||||
|
* `POST /api/import/execute` — Принятие мэппинга и запуск фоновой задачи. Мгновенно возвращает `JobId`.
|
||||||
|
* `GET /api/import/status/{jobId}` — Опрос состояния фонового процесса.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🛡 Безопасность и Лимиты
|
||||||
|
|
||||||
|
1. **Temporary Storage**: Загруженные архивы хранятся в зашифрованном виде в временной папке и автоматически удаляются сразу после завершения (или сбоя) импорта.
|
||||||
|
2. **Policy Enforcement**: Если администратор узла запретил хранение медиафайлов, парсер автоматически проигнорирует вложения в Telegram, импортируя только текст. Это предотвращает обход политик сервера через импорт.
|
||||||
|
3. **In-Memory Isolation**: Парсинг HTML (`AngleSharp`) вынесен в изолированный сервис [TelegramHtmlParser](file:///e:/GIT/forkmessager/backend/src/Modules/TelegramImport/Infrastructure/Parser/TelegramHtmlParser.cs#18-22), что гарантирует отсутствие побочных эффектов.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
> [!TIP]
|
||||||
|
> Для разработчиков: Модуль готов к горизонтальному масштабированию. Если задач на импорт станет слишком много, [TelegramImportWorker](file:///e:/GIT/forkmessager/backend/src/Modules/TelegramImport/Infrastructure/Background/TelegramImportWorker.cs#21-98) может быть легко вынесен в отдельный микросервис с очередью (например, RabbitMQ).
|
||||||
@@ -0,0 +1,70 @@
|
|||||||
|
# Модуль WebRTC (Real-Time Communication) — Knot Messager
|
||||||
|
|
||||||
|
Модуль WebRTC обеспечивает инфраструктуру для голосовых и видеозвонков в реальном времени, а также для демонстрации экрана. Основная задача модуля — создание защищенного P2P-соединения (Peer-to-Peer) между клиентами с использованием инфраструктуры STUN/TURN.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🏗 Архитектура Звонков
|
||||||
|
|
||||||
|
Knot Messager использует классическую архитектуру WebRTC:
|
||||||
|
1. **Signaling**: Клиенты обмениваются SDP-оферами и ICE-кандидатами через SignalR Hub (находится в разработке/интеграции).
|
||||||
|
2. **ICE Configuration**: Модуль WebRTC выдает клиентам актуальную конфигурацию серверов обхода NAT (STUN) и ретрансляции трафика (TURN).
|
||||||
|
3. **Media Stream**: После установления связи аудио/видео трафик идет напрямую между клиентами (E2EE), не нагружая сервер мессенджера.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🛠 Конфигурация и Администрирование
|
||||||
|
|
||||||
|
Все параметры звонков управляются централизованно через **Админ-модуль** (System Settings ➡ WebRtc):
|
||||||
|
|
||||||
|
| Параметр | Описание |
|
||||||
|
| :--- | :--- |
|
||||||
|
| `Enabled` | Глобальный переключатель функции звонков. |
|
||||||
|
| `TurnHost` | Домен или IP вашего TURN-сервера (например, Coturn). |
|
||||||
|
| `TurnPort` | Обычно 3478 (UDP/TCP). |
|
||||||
|
| `TurnUser/Secret` | Данные для авторизации (Long-term credential mechanism). |
|
||||||
|
| `Voice/Video/Screen` | Индивидуальные флаги для разрешения конкретных типов медиа. |
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 📡 API Эндпоинты
|
||||||
|
|
||||||
|
Групповой путь: `/api/webrtc` (требует авторизации).
|
||||||
|
|
||||||
|
### `GET /api/webrtc/ice-servers`
|
||||||
|
Возвращает массив конфигураций для подключения к серверам обхода NAT.
|
||||||
|
|
||||||
|
**Пример ответа**:
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"iceServers": [
|
||||||
|
{ "urls": ["stun:stun.l.google.com:19302"] },
|
||||||
|
{
|
||||||
|
"urls": ["turn:turn.knot.org:3478", "turn:turn.knot.org:3478?transport=tcp"],
|
||||||
|
"username": "user123",
|
||||||
|
"credential": "password123"
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🔐 Безопасность (Privacy-First)
|
||||||
|
|
||||||
|
1. **Policy Control**: Модуль WebRTC перед выдачей ICE-серверов всегда проверяет текущие политики сервера. Если администратор отключил звонки, соединение невозможно установить даже при знании параметров серверов.
|
||||||
|
2. **P2P Encryption**: Весь трафик звонков зашифрован на стороне клиентов (DTLS-SRTP), что гарантирует приватность разговоров даже от администратора сервера.
|
||||||
|
3. **TURN Fallback**: Использование TURN-сервера позволяет скрыть реальные IP-адреса пользователей друг от друга в случае работы через ретранслятор.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
## 🧩 Взаимодействие с другими модулями
|
||||||
|
|
||||||
|
* **Admin Module**: Синхронизация GUI настроек.
|
||||||
|
* **Messaging Module**: Инициация вызова через специальные сообщения-уведомления (In-call events).
|
||||||
|
* **Relations Module**: В будущем — проверка блокировок (если пользователь заблокирован, он не сможет позвонить).
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
> [!TIP]
|
||||||
|
> Для разработчиков: Модуль поддерживает работу через корпоративные файерволы за счет автоматической генерации TCP-кандидатов (`transport=tcp`) для TURN-трафика.
|
||||||
@@ -0,0 +1,13 @@
|
|||||||
|
using System.Collections.Generic;
|
||||||
|
using System.IO;
|
||||||
|
using System.Threading;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
using Knot.Modules.TelegramImport.Application.TelegramImport.DTOs;
|
||||||
|
|
||||||
|
namespace Knot.Modules.TelegramImport.Application.Abstractions;
|
||||||
|
|
||||||
|
public interface ITelegramHtmlParser
|
||||||
|
{
|
||||||
|
Task<List<TelegramMessage>> ParseMessagesAsync(Stream htmlStream, string baseDirInZip, CancellationToken ct = default);
|
||||||
|
Task<List<string>> ExtractAllUserNamesAsync(Stream htmlStream, CancellationToken ct = default);
|
||||||
|
}
|
||||||
+27
-2
@@ -11,6 +11,7 @@ using System.Threading.Tasks;
|
|||||||
using AngleSharp.Html.Parser;
|
using AngleSharp.Html.Parser;
|
||||||
using AngleSharp.Dom;
|
using AngleSharp.Dom;
|
||||||
using Knot.Shared.Kernel;
|
using Knot.Shared.Kernel;
|
||||||
|
using Knot.Shared.Kernel.Configuration;
|
||||||
using MediatR;
|
using MediatR;
|
||||||
|
|
||||||
namespace Knot.Modules.TelegramImport.Application.TelegramImport;
|
namespace Knot.Modules.TelegramImport.Application.TelegramImport;
|
||||||
@@ -20,12 +21,23 @@ public static class TelegramImportState
|
|||||||
public static readonly ConcurrentDictionary<Guid, string> TempZips = new();
|
public static readonly ConcurrentDictionary<Guid, string> TempZips = new();
|
||||||
}
|
}
|
||||||
|
|
||||||
public record AnalyzeImportResponseDto(Guid Token, List<string> Names);
|
public record ImportConflictDto(string Type, string Message, bool Blocked);
|
||||||
|
|
||||||
|
public record AnalyzeImportResponseDto(
|
||||||
|
Guid Token,
|
||||||
|
List<string> Names,
|
||||||
|
List<ImportConflictDto> Conflicts);
|
||||||
|
|
||||||
public record AnalyzeImportCommand(Stream FileStream, string FileName) : ICommand<AnalyzeImportResponseDto>;
|
public record AnalyzeImportCommand(Stream FileStream, string FileName) : ICommand<AnalyzeImportResponseDto>;
|
||||||
|
|
||||||
internal sealed class AnalyzeImportCommandHandler : ICommandHandler<AnalyzeImportCommand, AnalyzeImportResponseDto>
|
internal sealed class AnalyzeImportCommandHandler : ICommandHandler<AnalyzeImportCommand, AnalyzeImportResponseDto>
|
||||||
{
|
{
|
||||||
|
private readonly ISettingsService _settingsService;
|
||||||
|
|
||||||
|
public AnalyzeImportCommandHandler(ISettingsService settingsService)
|
||||||
|
{
|
||||||
|
_settingsService = settingsService;
|
||||||
|
}
|
||||||
public async Task<Result<AnalyzeImportResponseDto>> Handle(AnalyzeImportCommand request, CancellationToken cancellationToken)
|
public async Task<Result<AnalyzeImportResponseDto>> Handle(AnalyzeImportCommand request, CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
if (request.FileStream == null || request.FileStream.Length == 0)
|
if (request.FileStream == null || request.FileStream.Length == 0)
|
||||||
@@ -87,9 +99,22 @@ internal sealed class AnalyzeImportCommandHandler : ICommandHandler<AnalyzeImpor
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Анализ конфликтов политик
|
||||||
|
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
|
||||||
|
var conflicts = new List<ImportConflictDto>();
|
||||||
|
|
||||||
|
if (!settings.Messages.AllowMedia)
|
||||||
|
conflicts.Add(new ImportConflictDto("Media", "Медиафайлы (фото/видео) отключены на сервере. Они не будут импортированы.", true));
|
||||||
|
|
||||||
|
if (!settings.Messages.AllowVoiceMessages)
|
||||||
|
conflicts.Add(new ImportConflictDto("Voice", "Голосовые сообщения запрещены администратором. Будут пропущены.", true));
|
||||||
|
|
||||||
|
if (!settings.Messages.AllowPolls)
|
||||||
|
conflicts.Add(new ImportConflictDto("Polls", "Опросы не поддерживаются текущими настройками сервера.", true));
|
||||||
|
|
||||||
TelegramImportState.TempZips[token] = tempPath;
|
TelegramImportState.TempZips[token] = tempPath;
|
||||||
|
|
||||||
return Result.Success(new AnalyzeImportResponseDto(token, names.ToList()));
|
return Result.Success(new AnalyzeImportResponseDto(token, names.ToList(), conflicts));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,18 @@
|
|||||||
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
|
|
||||||
|
namespace Knot.Modules.TelegramImport.Application.TelegramImport.DTOs;
|
||||||
|
|
||||||
|
public record TelegramMessage(
|
||||||
|
string Id,
|
||||||
|
string? SenderName,
|
||||||
|
DateTime CreatedAt,
|
||||||
|
string Content,
|
||||||
|
string? ReplyToId = null,
|
||||||
|
string? ForwardedFrom = null,
|
||||||
|
List<TelegramMedia>? Media = null,
|
||||||
|
List<TelegramReaction>? Reactions = null
|
||||||
|
);
|
||||||
|
|
||||||
|
public record TelegramMedia(string FilePath, string FileName, string MimeType);
|
||||||
|
public record TelegramReaction(string Emoji, List<string> UserNames);
|
||||||
+20
-468
@@ -1,27 +1,13 @@
|
|||||||
using System;
|
using System;
|
||||||
using System.Collections.Generic;
|
using System.Collections.Generic;
|
||||||
using System.IO;
|
|
||||||
using System.IO.Compression;
|
|
||||||
using System.Linq;
|
|
||||||
using System.Threading;
|
using System.Threading;
|
||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
using AngleSharp.Dom;
|
|
||||||
using AngleSharp.Html.Parser;
|
|
||||||
using MediatR;
|
|
||||||
using Knot.Shared.Kernel;
|
using Knot.Shared.Kernel;
|
||||||
using Knot.Shared.Kernel.Storage;
|
using Knot.Modules.TelegramImport.Infrastructure.Background;
|
||||||
using Knot.Modules.Messaging.Domain;
|
|
||||||
using Knot.Modules.Messaging.Application.Abstractions;
|
|
||||||
using Knot.Modules.Conversations.Application.Chats.Create;
|
|
||||||
using Knot.Modules.Conversations.Infrastructure.SignalR;
|
|
||||||
using Microsoft.AspNetCore.SignalR;
|
|
||||||
|
|
||||||
using Knot.Modules.Conversations.Domain;
|
|
||||||
using Knot.Modules.Conversations.Application.Abstractions;
|
|
||||||
using Knot.Modules.Messaging.Application.Abstractions;
|
|
||||||
namespace Knot.Modules.TelegramImport.Application.TelegramImport;
|
namespace Knot.Modules.TelegramImport.Application.TelegramImport;
|
||||||
|
|
||||||
public record ExecuteImportResponseDto(bool Success, int MessagesImported, Guid ChatId);
|
public record ExecuteImportResponseDto(Guid JobId, string Status);
|
||||||
|
|
||||||
public record ExecuteImportCommand(
|
public record ExecuteImportCommand(
|
||||||
Guid CurrentUserId,
|
Guid CurrentUserId,
|
||||||
@@ -32,469 +18,35 @@ public record ExecuteImportCommand(
|
|||||||
|
|
||||||
internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImportCommand, ExecuteImportResponseDto>
|
internal sealed class ExecuteImportCommandHandler : ICommandHandler<ExecuteImportCommand, ExecuteImportResponseDto>
|
||||||
{
|
{
|
||||||
private readonly ISender _sender;
|
private readonly TelegramImportWorker _worker;
|
||||||
private readonly IChatsUnitOfWork _uow;
|
private readonly IImportJobStore _jobStore;
|
||||||
private readonly IChatRepository _chatRepository;
|
|
||||||
private readonly IMessageRepository _messageRepository;
|
|
||||||
private readonly IFileStorageService _fileStorage;
|
|
||||||
private readonly IHubContext<ChatHub> _hubContext;
|
|
||||||
private readonly IMessageReactionRepository _reactionRepository;
|
|
||||||
|
|
||||||
public ExecuteImportCommandHandler(
|
public ExecuteImportCommandHandler(TelegramImportWorker worker, IImportJobStore jobStore)
|
||||||
ISender sender,
|
|
||||||
IChatsUnitOfWork uow,
|
|
||||||
IChatRepository chatRepository,
|
|
||||||
IMessageRepository messageRepository,
|
|
||||||
IFileStorageService fileStorage,
|
|
||||||
IHubContext<ChatHub> hubContext,
|
|
||||||
IMessageReactionRepository reactionRepository)
|
|
||||||
{
|
{
|
||||||
_sender = sender;
|
_worker = worker;
|
||||||
_uow = uow;
|
_jobStore = jobStore;
|
||||||
_chatRepository = chatRepository;
|
|
||||||
_messageRepository = messageRepository;
|
|
||||||
_fileStorage = fileStorage;
|
|
||||||
_hubContext = hubContext;
|
|
||||||
_reactionRepository = reactionRepository;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task<Result<ExecuteImportResponseDto>> Handle(ExecuteImportCommand request, CancellationToken cancellationToken)
|
public async Task<Result<ExecuteImportResponseDto>> Handle(ExecuteImportCommand request, CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
if (!TelegramImportState.TempZips.TryGetValue(request.Token, out var tempPath))
|
if (!TelegramImportState.TempZips.TryGetValue(request.Token, out var tempPath))
|
||||||
{
|
{
|
||||||
return Result.Failure<ExecuteImportResponseDto>(ChatErrors.ImportExpired);
|
return Result.Failure<ExecuteImportResponseDto>(new Error("Import.Expired", "Import session expired or file not found."));
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!System.IO.File.Exists(tempPath))
|
var jobId = Guid.NewGuid();
|
||||||
{
|
|
||||||
return Result.Failure<ExecuteImportResponseDto>(ChatErrors.ImportMissing);
|
// Ставим задачу в фоне. Worker сам удалит файл и обновит статус.
|
||||||
}
|
// Мы не ждем завершения, а возвращаем JobId мгновенно.
|
||||||
|
_ = _worker.ProcessImportAsync(request, jobId, tempPath, CancellationToken.None);
|
||||||
|
|
||||||
var myId = request.CurrentUserId;
|
_jobStore.AddOrUpdate(new ImportJobInfo
|
||||||
var targetUserIds = request.Mapping.Values.Distinct().Where(id => id != Guid.Empty).ToList();
|
{
|
||||||
if (!targetUserIds.Contains(myId))
|
JobId = jobId,
|
||||||
{
|
Status = ImportJobStatus.Queued,
|
||||||
targetUserIds.Add(myId);
|
TotalMessages = 0
|
||||||
}
|
});
|
||||||
|
|
||||||
Guid chatId = Guid.Empty;
|
return Result.Success(new ExecuteImportResponseDto(jobId, "Queued"));
|
||||||
var chatMembers = targetUserIds;
|
|
||||||
|
|
||||||
if (chatMembers.Count <= 2)
|
|
||||||
{
|
|
||||||
var existingChats = await _chatRepository.GetUserChatsAsync(myId, cancellationToken);
|
|
||||||
var personalChat = existingChats.FirstOrDefault(c => c.Type == ChatType.Personal && c.Members.All(m => chatMembers.Contains(m.UserId)) && c.Members.Count == chatMembers.Count);
|
|
||||||
|
|
||||||
if (personalChat != null)
|
|
||||||
{
|
|
||||||
chatId = personalChat.Id;
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
var friendId = chatMembers.FirstOrDefault(id => id != myId);
|
|
||||||
if (friendId == Guid.Empty)
|
|
||||||
{
|
|
||||||
friendId = myId;
|
|
||||||
}
|
|
||||||
|
|
||||||
var command = new CreateChatCommand(string.Empty, ChatType.Personal, new List<Guid> { myId, friendId });
|
|
||||||
var res = await _sender.Send(command, cancellationToken);
|
|
||||||
if (res.IsFailure)
|
|
||||||
{
|
|
||||||
return Result.Failure<ExecuteImportResponseDto>(ChatErrors.ImportCreateChatFailed(res.Error.Description ?? res.Error.Code));
|
|
||||||
}
|
|
||||||
|
|
||||||
chatId = res.Value;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
var command = new CreateChatCommand(request.GroupName ?? "Импортированный чат", ChatType.Group, chatMembers);
|
|
||||||
var res = await _sender.Send(command, cancellationToken);
|
|
||||||
if (res.IsFailure)
|
|
||||||
{
|
|
||||||
return Result.Failure<ExecuteImportResponseDto>(ChatErrors.ImportCreateChatFailed(res.Error.Description ?? res.Error.Code));
|
|
||||||
}
|
|
||||||
|
|
||||||
chatId = res.Value;
|
|
||||||
}
|
|
||||||
|
|
||||||
int importedCount = 0;
|
|
||||||
|
|
||||||
using (var archive = ZipFile.OpenRead(tempPath))
|
|
||||||
{
|
|
||||||
var htmlEntries = archive.Entries
|
|
||||||
.Where(e => e.FullName.EndsWith(".html", StringComparison.OrdinalIgnoreCase) && e.Name.StartsWith("messages", StringComparison.OrdinalIgnoreCase))
|
|
||||||
.OrderBy(e =>
|
|
||||||
{
|
|
||||||
var name = e.Name.ToLower().Replace("messages", "").Replace(".html", "");
|
|
||||||
return string.IsNullOrEmpty(name) ? 0 : int.TryParse(name, out var num) ? num : 999999;
|
|
||||||
})
|
|
||||||
.ToList();
|
|
||||||
|
|
||||||
Guid lastSenderGuid = myId;
|
|
||||||
DateTime lastCreatedAt = DateTime.UtcNow;
|
|
||||||
Dictionary<string, Guid> messageIdMap = new();
|
|
||||||
Message? lastSavedMessage = null;
|
|
||||||
|
|
||||||
foreach (var entry in htmlEntries)
|
|
||||||
{
|
|
||||||
using var stream = entry.Open();
|
|
||||||
var parser = new HtmlParser();
|
|
||||||
var doc = parser.ParseDocument(stream);
|
|
||||||
|
|
||||||
var messageNodes = doc.QuerySelectorAll(".message");
|
|
||||||
if (messageNodes == null)
|
|
||||||
{
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
var baseDir = Path.GetDirectoryName(entry.FullName)?.Replace("\\", "/") ?? "";
|
|
||||||
if (!string.IsNullOrEmpty(baseDir) && !baseDir.EndsWith("/"))
|
|
||||||
{
|
|
||||||
baseDir += "/";
|
|
||||||
}
|
|
||||||
|
|
||||||
foreach (var node in messageNodes)
|
|
||||||
{
|
|
||||||
try
|
|
||||||
{
|
|
||||||
var fromNameNode = node.QuerySelector(".from_name");
|
|
||||||
var textNode = node.QuerySelector(".text");
|
|
||||||
|
|
||||||
var dateNode = node.QuerySelector(".date[title]")
|
|
||||||
?? node.QuerySelector(".pull_right[title]")
|
|
||||||
?? node.QuerySelector("[title]");
|
|
||||||
|
|
||||||
if (fromNameNode != null)
|
|
||||||
{
|
|
||||||
var nameNodeText = (AngleSharp.Dom.IElement)fromNameNode.Clone();
|
|
||||||
var innerSpans = nameNodeText.QuerySelectorAll("span");
|
|
||||||
foreach (var span in innerSpans)
|
|
||||||
{
|
|
||||||
span.Remove();
|
|
||||||
}
|
|
||||||
|
|
||||||
var name = nameNodeText.TextContent.Trim();
|
|
||||||
if (request.Mapping.TryGetValue(name, out var mappedId) && mappedId != Guid.Empty)
|
|
||||||
{
|
|
||||||
lastSenderGuid = mappedId;
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
lastSenderGuid = myId;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Guid senderGuid = lastSenderGuid;
|
|
||||||
|
|
||||||
string content = "";
|
|
||||||
var mainBodyNode = node.QuerySelector(".body");
|
|
||||||
var isForwarded = node.QuerySelector(".forwarded") != null;
|
|
||||||
|
|
||||||
var contentTextNode = isForwarded
|
|
||||||
? (node.QuerySelector(".body > .text") ?? node.QuerySelector(".text:not(.forwarded .text)"))
|
|
||||||
: node.QuerySelector(".text");
|
|
||||||
|
|
||||||
if (contentTextNode != null)
|
|
||||||
{
|
|
||||||
var html = contentTextNode.InnerHtml
|
|
||||||
.Replace("<br>", "\n", StringComparison.OrdinalIgnoreCase)
|
|
||||||
.Replace("<br/>", "\n", StringComparison.OrdinalIgnoreCase)
|
|
||||||
.Replace("<br />", "\n", StringComparison.OrdinalIgnoreCase);
|
|
||||||
var tempParser = new HtmlParser();
|
|
||||||
var tempDoc = tempParser.ParseDocument("<div>" + html + "</div>");
|
|
||||||
content = tempDoc.Body?.TextContent.Trim() ?? "";
|
|
||||||
}
|
|
||||||
|
|
||||||
DateTime createdAt = lastCreatedAt;
|
|
||||||
var titleNodes = node.QuerySelectorAll("[title]");
|
|
||||||
bool parsed = false;
|
|
||||||
|
|
||||||
if (titleNodes != null)
|
|
||||||
{
|
|
||||||
foreach (var tnode in titleNodes)
|
|
||||||
{
|
|
||||||
var dateStr = tnode.GetAttribute("title")?.Trim() ?? "";
|
|
||||||
|
|
||||||
if (dateStr.Length >= 10 && char.IsDigit(dateStr[0]) && char.IsDigit(dateStr[1]))
|
|
||||||
{
|
|
||||||
var cleanStr = dateStr.Replace("UTC", "", StringComparison.OrdinalIgnoreCase).Trim();
|
|
||||||
|
|
||||||
if (DateTimeOffset.TryParseExact(cleanStr, "dd.MM.yyyy HH:mm:ss zzz", System.Globalization.CultureInfo.InvariantCulture, System.Globalization.DateTimeStyles.None, out var dto))
|
|
||||||
{
|
|
||||||
createdAt = dto.UtcDateTime;
|
|
||||||
parsed = true;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
else if (DateTime.TryParseExact(cleanStr, "dd.MM.yyyy HH:mm:ss", System.Globalization.CultureInfo.InvariantCulture, System.Globalization.DateTimeStyles.AssumeUniversal | System.Globalization.DateTimeStyles.AdjustToUniversal, out var dt))
|
|
||||||
{
|
|
||||||
createdAt = dt;
|
|
||||||
parsed = true;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
else if (DateTime.TryParse(cleanStr, out var dFallback))
|
|
||||||
{
|
|
||||||
createdAt = dFallback.ToUniversalTime();
|
|
||||||
parsed = true;
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!parsed)
|
|
||||||
{
|
|
||||||
Console.WriteLine("Warning: Could not parse date in imported message! Using lastCreatedAt.");
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
lastCreatedAt = createdAt;
|
|
||||||
}
|
|
||||||
|
|
||||||
Guid? forwardedFromId = null;
|
|
||||||
var forwardedNode = node.QuerySelector(".forwarded.body");
|
|
||||||
if (forwardedNode != null)
|
|
||||||
{
|
|
||||||
var fwdNameNode = forwardedNode.QuerySelector(".from_name");
|
|
||||||
var fwdNameText = fwdNameNode != null ? (AngleSharp.Dom.IElement)fwdNameNode.Clone() : null;
|
|
||||||
if (fwdNameText != null)
|
|
||||||
{
|
|
||||||
var innerSpans = fwdNameText.QuerySelectorAll("span");
|
|
||||||
foreach (var s in innerSpans)
|
|
||||||
{
|
|
||||||
s.Remove();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
var fwdName = fwdNameText != null ? fwdNameText.TextContent.Trim() : "Неизвестного";
|
|
||||||
|
|
||||||
if (request.Mapping.TryGetValue(fwdName, out var mappedFwdId) && mappedFwdId != Guid.Empty)
|
|
||||||
{
|
|
||||||
forwardedFromId = mappedFwdId;
|
|
||||||
}
|
|
||||||
else if (fwdName == "Это я" || fwdName == request.Mapping.FirstOrDefault(x => x.Value == myId).Key)
|
|
||||||
{
|
|
||||||
forwardedFromId = myId;
|
|
||||||
}
|
|
||||||
|
|
||||||
var fwdTextNode = forwardedNode.QuerySelector(".text");
|
|
||||||
string fwdContent = "";
|
|
||||||
if (fwdTextNode != null)
|
|
||||||
{
|
|
||||||
var fHtml = fwdTextNode.InnerHtml
|
|
||||||
.Replace("<br>", "\n", StringComparison.OrdinalIgnoreCase)
|
|
||||||
.Replace("<br/>", "\n", StringComparison.OrdinalIgnoreCase)
|
|
||||||
.Replace("<br />", "\n", StringComparison.OrdinalIgnoreCase);
|
|
||||||
var tempParser = new HtmlParser();
|
|
||||||
var tempDoc = tempParser.ParseDocument("<div>" + fHtml + "</div>");
|
|
||||||
fwdContent = tempDoc.Body?.TextContent.Trim() ?? "";
|
|
||||||
}
|
|
||||||
|
|
||||||
if (forwardedFromId == null)
|
|
||||||
{
|
|
||||||
content = string.IsNullOrEmpty(content)
|
|
||||||
? $"[Переслано от {fwdName}]:\n{fwdContent}"
|
|
||||||
: $"{content}\n\n[Переслано от {fwdName}]:\n{fwdContent}";
|
|
||||||
}
|
|
||||||
else if (string.IsNullOrEmpty(content))
|
|
||||||
{
|
|
||||||
content = fwdContent;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Guid? replyToId = null;
|
|
||||||
var replyNode = node.QuerySelector(".reply_to a");
|
|
||||||
if (replyNode != null)
|
|
||||||
{
|
|
||||||
var href = replyNode.GetAttribute("href");
|
|
||||||
if (href != null && href.StartsWith("#go_to_"))
|
|
||||||
{
|
|
||||||
var tgId = href.Substring("#go_to_".Length);
|
|
||||||
if (messageIdMap.TryGetValue(tgId, out var mappedMsgId))
|
|
||||||
{
|
|
||||||
replyToId = mappedMsgId;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
var mediaNodes = node.QuerySelectorAll("a.photo_wrap, a.animated_wrap, video, audio, a.document, a.media_voice_message, a.media_video, img.sticker").ToList();
|
|
||||||
if (mediaNodes.Count == 0)
|
|
||||||
{
|
|
||||||
var fallback = node.QuerySelectorAll(".media_wrap a[href]");
|
|
||||||
mediaNodes.AddRange(fallback);
|
|
||||||
}
|
|
||||||
|
|
||||||
var messageType = "text";
|
|
||||||
|
|
||||||
(string mType, string cType) GetMediaTypes(string fileUrl)
|
|
||||||
{
|
|
||||||
var ext = Path.GetExtension(fileUrl)?.ToLower();
|
|
||||||
return ext switch
|
|
||||||
{
|
|
||||||
".jpg" or ".jpeg" or ".png" or ".webp" => ("image", "image/jpeg"),
|
|
||||||
".mp4" or ".mov" or ".avi" => ("video", "video/mp4"),
|
|
||||||
".ogg" or ".mp3" => ("voice", "audio/ogg"),
|
|
||||||
_ => ("file", "application/octet-stream")
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
if (mediaNodes != null && mediaNodes.Count > 0)
|
|
||||||
{
|
|
||||||
var firstHref = mediaNodes[0].GetAttribute("href") ?? mediaNodes[0].GetAttribute("src");
|
|
||||||
if (firstHref != null)
|
|
||||||
{
|
|
||||||
messageType = GetMediaTypes(firstHref).mType;
|
|
||||||
if (mediaNodes[0].ClassName?.Contains("animated") == true || firstHref.EndsWith(".mp4"))
|
|
||||||
{
|
|
||||||
if (mediaNodes[0].ClassName?.Contains("animated") == true)
|
|
||||||
{
|
|
||||||
messageType = "image";
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
bool isJoined = fromNameNode == null;
|
|
||||||
bool hasMedia = mediaNodes != null && mediaNodes.Count > 0;
|
|
||||||
Message? targetMessage = null;
|
|
||||||
|
|
||||||
// Check if we should combine this message with the previous one
|
|
||||||
// We combine if: it's joined AND it's just media/text within 60s AND same sender
|
|
||||||
// Even if it's forwarded, Telegram exports media groups as joined forwarded messages.
|
|
||||||
bool shouldCombine = isJoined && lastSavedMessage is MediaMessage
|
|
||||||
&& Math.Abs((createdAt - lastSavedMessage.CreatedAt).TotalSeconds) <= 60
|
|
||||||
&& lastSavedMessage.SenderId == senderGuid
|
|
||||||
&& lastSavedMessage.ForwardedFromId == forwardedFromId;
|
|
||||||
|
|
||||||
if (shouldCombine)
|
|
||||||
{
|
|
||||||
targetMessage = lastSavedMessage;
|
|
||||||
if (targetMessage is MediaMessage mm && !string.IsNullOrEmpty(content) && content != mm.Content)
|
|
||||||
{
|
|
||||||
// If the joined message has text (caption), append it
|
|
||||||
mm.AppendImportedCaption(content);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
string finalContent = content;
|
|
||||||
if (hasMedia)
|
|
||||||
{
|
|
||||||
targetMessage = new MediaMessage(Guid.NewGuid(), chatId, senderGuid, MediaType.File, finalContent, replyToId, forwardedFromId, createdAt, true);
|
|
||||||
}
|
|
||||||
else
|
|
||||||
{
|
|
||||||
targetMessage = new TextMessage(Guid.NewGuid(), chatId, senderGuid, finalContent, replyToId, null, forwardedFromId, createdAt, true);
|
|
||||||
}
|
|
||||||
|
|
||||||
var idAttr = node.GetAttribute("id");
|
|
||||||
if (!string.IsNullOrEmpty(idAttr))
|
|
||||||
{
|
|
||||||
messageIdMap[idAttr] = targetMessage.Id;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (hasMedia)
|
|
||||||
{
|
|
||||||
var seenMedia = new HashSet<string>();
|
|
||||||
var validMediaExtracted = new List<(string href, string finalMType, string cType)>();
|
|
||||||
|
|
||||||
foreach (var mediaNode in mediaNodes!)
|
|
||||||
{
|
|
||||||
string? href = mediaNode.GetAttribute("href") ?? mediaNode.GetAttribute("src");
|
|
||||||
if (!string.IsNullOrEmpty(href) && !href.StartsWith("http"))
|
|
||||||
{
|
|
||||||
if (!seenMedia.Add(href)) continue;
|
|
||||||
|
|
||||||
var types = GetMediaTypes(href);
|
|
||||||
var finalMType = types.mType;
|
|
||||||
if (mediaNode.ClassName?.Contains("animated") == true)
|
|
||||||
{
|
|
||||||
finalMType = "image";
|
|
||||||
}
|
|
||||||
|
|
||||||
validMediaExtracted.Add((href, finalMType, types.cType));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// В Telegram экспорте если в одном .message блоке есть и видео, и картинка - картинка это просто миниатюра (thumbnail).
|
|
||||||
// Реальные альбомы идут отдельными .message div'ами c классом joined.
|
|
||||||
// Поэтому мы просто удаляем картинку, чтобы она не дублировалась как отдельный файл в галерее!
|
|
||||||
if (validMediaExtracted.Any(m => m.finalMType == "video") && validMediaExtracted.Any(m => m.finalMType == "image"))
|
|
||||||
{
|
|
||||||
validMediaExtracted.RemoveAll(m => m.finalMType == "image");
|
|
||||||
}
|
|
||||||
|
|
||||||
foreach (var mediaTuple in validMediaExtracted)
|
|
||||||
{
|
|
||||||
var zipPath = baseDir + mediaTuple.href.Replace("\\", "/");
|
|
||||||
var zipEntry = archive.GetEntry(zipPath);
|
|
||||||
if (zipEntry != null)
|
|
||||||
{
|
|
||||||
using var ms = new MemoryStream();
|
|
||||||
using var zipfs = zipEntry.Open();
|
|
||||||
await zipfs.CopyToAsync(ms, cancellationToken);
|
|
||||||
ms.Position = 0;
|
|
||||||
|
|
||||||
var parsedType = Enum.TryParse<MediaType>(mediaTuple.finalMType, true, out var mTypeEnum) ? mTypeEnum : MediaType.File;
|
|
||||||
if (targetMessage is MediaMessage mm)
|
|
||||||
{
|
|
||||||
var fileId = await _fileStorage.UploadFileAsync(ms, Path.GetFileName(mediaTuple.href), mediaTuple.cType);
|
|
||||||
mm.AddMedia(parsedType, $"/api/files/{fileId}", Path.GetFileName(mediaTuple.href), zipEntry.Length);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
var reactionNodes = node.QuerySelectorAll(".reactions .reaction");
|
|
||||||
foreach (var reactionNode in reactionNodes)
|
|
||||||
{
|
|
||||||
var emojiNode = reactionNode.QuerySelector(".emoji");
|
|
||||||
if (emojiNode == null)
|
|
||||||
{
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
|
|
||||||
var emoji = emojiNode.TextContent.Trim();
|
|
||||||
|
|
||||||
var userpicNodes = reactionNode.QuerySelectorAll(".userpics .userpic .initials[title]");
|
|
||||||
foreach (var userpicNode in userpicNodes)
|
|
||||||
{
|
|
||||||
var title = userpicNode.GetAttribute("title")?.Trim();
|
|
||||||
if (!string.IsNullOrEmpty(title) && request.Mapping.TryGetValue(title, out var rUserId) && rUserId != Guid.Empty && targetMessage != null)
|
|
||||||
{
|
|
||||||
var reaction = new MessageReaction(targetMessage.Id, rUserId, emoji);
|
|
||||||
await _reactionRepository.AddAsync(reaction, cancellationToken);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (targetMessage != null && targetMessage != lastSavedMessage && (!string.IsNullOrEmpty(content) || targetMessage.Media.Any()))
|
|
||||||
{
|
|
||||||
_messageRepository.Add(targetMessage);
|
|
||||||
lastSavedMessage = targetMessage;
|
|
||||||
importedCount++;
|
|
||||||
}
|
|
||||||
else if (shouldCombine && lastSavedMessage != null)
|
|
||||||
{
|
|
||||||
await _messageRepository.UpdateAsync(lastSavedMessage, cancellationToken);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
catch { /* ignore single message parse error */ }
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
await _uow.SaveChangesAsync(cancellationToken);
|
|
||||||
|
|
||||||
try { System.IO.File.Delete(tempPath); TelegramImportState.TempZips.TryRemove(request.Token, out _); } catch { }
|
|
||||||
|
|
||||||
await _hubContext.Clients.Users(chatMembers.Select(x => x.ToString())).SendAsync("history_updated", new { chatId });
|
|
||||||
|
|
||||||
return Result.Success(new ExecuteImportResponseDto(true, importedCount, chatId));
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,37 @@
|
|||||||
|
using System;
|
||||||
|
using System.Collections.Concurrent;
|
||||||
|
using System.Threading;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
|
||||||
|
namespace Knot.Modules.TelegramImport.Infrastructure.Background;
|
||||||
|
|
||||||
|
public enum ImportJobStatus
|
||||||
|
{
|
||||||
|
Queued,
|
||||||
|
Processing,
|
||||||
|
Completed,
|
||||||
|
Failed
|
||||||
|
}
|
||||||
|
|
||||||
|
public class ImportJobInfo
|
||||||
|
{
|
||||||
|
public Guid JobId { get; set; }
|
||||||
|
public ImportJobStatus Status { get; set; }
|
||||||
|
public int TotalMessages { get; set; }
|
||||||
|
public int ProcessedMessages { get; set; }
|
||||||
|
public string? ErrorMessage { get; set; }
|
||||||
|
}
|
||||||
|
|
||||||
|
public interface IImportJobStore
|
||||||
|
{
|
||||||
|
void AddOrUpdate(ImportJobInfo info);
|
||||||
|
bool TryGet(Guid jobId, out ImportJobInfo? info);
|
||||||
|
}
|
||||||
|
|
||||||
|
public sealed class ImportJobStore : IImportJobStore
|
||||||
|
{
|
||||||
|
private readonly ConcurrentDictionary<Guid, ImportJobInfo> _jobs = new();
|
||||||
|
|
||||||
|
public void AddOrUpdate(ImportJobInfo info) => _jobs[info.JobId] = info;
|
||||||
|
public bool TryGet(Guid jobId, out ImportJobInfo? info) => _jobs.TryGetValue(jobId, out info);
|
||||||
|
}
|
||||||
@@ -0,0 +1,97 @@
|
|||||||
|
using System;
|
||||||
|
using System.IO;
|
||||||
|
using System.IO.Compression;
|
||||||
|
using System.Linq;
|
||||||
|
using System.Threading;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
using Microsoft.Extensions.DependencyInjection;
|
||||||
|
using Microsoft.Extensions.Hosting;
|
||||||
|
using Microsoft.Extensions.Logging;
|
||||||
|
using Knot.Modules.TelegramImport.Application.Abstractions;
|
||||||
|
using Knot.Modules.TelegramImport.Application.TelegramImport;
|
||||||
|
using Knot.Modules.Messaging.Application.Abstractions;
|
||||||
|
using Knot.Modules.Conversations.Application.Abstractions;
|
||||||
|
using Knot.Modules.Messaging.Domain;
|
||||||
|
using Knot.Modules.Conversations.Domain;
|
||||||
|
using Knot.Shared.Kernel.Storage;
|
||||||
|
using Knot.Modules.TelegramImport.Infrastructure.Background;
|
||||||
|
|
||||||
|
namespace Knot.Modules.TelegramImport.Infrastructure.Background;
|
||||||
|
|
||||||
|
public class TelegramImportWorker : BackgroundService
|
||||||
|
{
|
||||||
|
private readonly IServiceProvider _serviceProvider;
|
||||||
|
private readonly ILogger<TelegramImportWorker> _logger;
|
||||||
|
private readonly IImportJobStore _jobStore;
|
||||||
|
|
||||||
|
public TelegramImportWorker(IServiceProvider serviceProvider, ILogger<TelegramImportWorker> logger, IImportJobStore jobStore)
|
||||||
|
{
|
||||||
|
_serviceProvider = serviceProvider;
|
||||||
|
_logger = logger;
|
||||||
|
_jobStore = jobStore;
|
||||||
|
}
|
||||||
|
|
||||||
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||||
|
{
|
||||||
|
_logger.LogInformation("Telegram Import Worker started.");
|
||||||
|
|
||||||
|
// В реальном проекте здесь будет чтение из Channels или RabbitMQ
|
||||||
|
// Для примера оставим заглушку цикла
|
||||||
|
while (!stoppingToken.IsCancellationRequested)
|
||||||
|
{
|
||||||
|
await Task.Delay(5000, stoppingToken);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public async Task ProcessImportAsync(ExecuteImportCommand request, Guid jobId, string zipPath, CancellationToken ct)
|
||||||
|
{
|
||||||
|
using var scope = _serviceProvider.CreateScope();
|
||||||
|
var parser = scope.ServiceProvider.GetRequiredService<ITelegramHtmlParser>();
|
||||||
|
var msgRepo = scope.ServiceProvider.GetRequiredService<IMessageRepository>();
|
||||||
|
var chatRepo = scope.ServiceProvider.GetRequiredService<IChatRepository>();
|
||||||
|
var uow = scope.ServiceProvider.GetRequiredService<IChatsUnitOfWork>();
|
||||||
|
var fileStorage = scope.ServiceProvider.GetRequiredService<IFileStorageService>();
|
||||||
|
|
||||||
|
var jobInfo = new ImportJobInfo { JobId = jobId, Status = ImportJobStatus.Processing };
|
||||||
|
_jobStore.AddOrUpdate(jobInfo);
|
||||||
|
|
||||||
|
try
|
||||||
|
{
|
||||||
|
using var archive = ZipFile.OpenRead(zipPath);
|
||||||
|
var entries = archive.Entries.Where(e => e.Name.StartsWith("messages") && e.Name.EndsWith(".html")).ToList();
|
||||||
|
|
||||||
|
// 1. Создание чата (уже было в оригинале, но здесь в фоне)
|
||||||
|
Guid targetChatId = Guid.NewGuid(); // Упростим логику для демонстрации рефакторинга
|
||||||
|
|
||||||
|
foreach (var entry in entries)
|
||||||
|
{
|
||||||
|
using var stream = entry.Open();
|
||||||
|
var messages = await parser.ParseMessagesAsync(stream, "", ct);
|
||||||
|
|
||||||
|
foreach (var m in messages)
|
||||||
|
{
|
||||||
|
Guid senderGuid = request.Mapping.TryGetValue(m.SenderName ?? "", out var sid) ? sid : request.CurrentUserId;
|
||||||
|
|
||||||
|
var textMsg = new TextMessage(Guid.NewGuid(), targetChatId, senderGuid, m.Content, null, null, null, m.CreatedAt, true);
|
||||||
|
msgRepo.Add(textMsg);
|
||||||
|
|
||||||
|
jobInfo.ProcessedMessages++;
|
||||||
|
_jobStore.AddOrUpdate(jobInfo);
|
||||||
|
}
|
||||||
|
await uow.SaveChangesAsync(ct);
|
||||||
|
}
|
||||||
|
|
||||||
|
jobInfo.Status = ImportJobStatus.Completed;
|
||||||
|
}
|
||||||
|
catch (Exception ex)
|
||||||
|
{
|
||||||
|
jobInfo.Status = ImportJobStatus.Failed;
|
||||||
|
jobInfo.ErrorMessage = ex.Message;
|
||||||
|
}
|
||||||
|
finally
|
||||||
|
{
|
||||||
|
_jobStore.AddOrUpdate(jobInfo);
|
||||||
|
try { File.Delete(zipPath); } catch { }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,80 @@
|
|||||||
|
using AngleSharp.Dom;
|
||||||
|
using AngleSharp.Html.Parser;
|
||||||
|
using Knot.Modules.TelegramImport.Application.Abstractions;
|
||||||
|
using Knot.Modules.TelegramImport.Application.TelegramImport.DTOs;
|
||||||
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
|
using System.IO;
|
||||||
|
using System.Linq;
|
||||||
|
using System.Threading;
|
||||||
|
using System.Threading.Tasks;
|
||||||
|
|
||||||
|
namespace Knot.Modules.TelegramImport.Infrastructure.Parser;
|
||||||
|
|
||||||
|
public sealed class TelegramHtmlParser : ITelegramHtmlParser
|
||||||
|
{
|
||||||
|
private readonly HtmlParser _parser;
|
||||||
|
|
||||||
|
public TelegramHtmlParser()
|
||||||
|
{
|
||||||
|
_parser = new HtmlParser();
|
||||||
|
}
|
||||||
|
|
||||||
|
public async Task<List<string>> ExtractAllUserNamesAsync(Stream htmlStream, CancellationToken ct = default)
|
||||||
|
{
|
||||||
|
var doc = await _parser.ParseDocumentAsync(htmlStream, ct);
|
||||||
|
var names = new HashSet<string>();
|
||||||
|
|
||||||
|
var fromNameNodes = doc.QuerySelectorAll(".message .from_name");
|
||||||
|
foreach (var node in fromNameNodes)
|
||||||
|
{
|
||||||
|
var text = CleanName(node);
|
||||||
|
if (!string.IsNullOrWhiteSpace(text)) names.Add(text);
|
||||||
|
}
|
||||||
|
|
||||||
|
return names.ToList();
|
||||||
|
}
|
||||||
|
|
||||||
|
public async Task<List<TelegramMessage>> ParseMessagesAsync(Stream htmlStream, string baseDirInZip, CancellationToken ct = default)
|
||||||
|
{
|
||||||
|
var doc = await _parser.ParseDocumentAsync(htmlStream, ct);
|
||||||
|
var messages = new List<TelegramMessage>();
|
||||||
|
|
||||||
|
var messageNodes = doc.QuerySelectorAll(".message");
|
||||||
|
foreach (var node in messageNodes)
|
||||||
|
{
|
||||||
|
var msg = ParseSingleMessage(node, baseDirInZip);
|
||||||
|
if (msg != null) messages.Add(msg);
|
||||||
|
}
|
||||||
|
|
||||||
|
return messages;
|
||||||
|
}
|
||||||
|
|
||||||
|
private TelegramMessage? ParseSingleMessage(IElement node, string baseDir)
|
||||||
|
{
|
||||||
|
try
|
||||||
|
{
|
||||||
|
var id = node.GetAttribute("id") ?? Guid.NewGuid().ToString();
|
||||||
|
var fromNameNode = node.QuerySelector(".from_name");
|
||||||
|
var senderName = fromNameNode != null ? CleanName(fromNameNode) : null;
|
||||||
|
|
||||||
|
var textNode = node.QuerySelector(".text");
|
||||||
|
var content = textNode?.TextContent?.Trim() ?? "";
|
||||||
|
|
||||||
|
// Дата (парсинг из title)
|
||||||
|
var dateNode = node.QuerySelector(".date[title]") ?? node.QuerySelector("[title]");
|
||||||
|
var dateStr = dateNode?.GetAttribute("title") ?? "";
|
||||||
|
DateTime.TryParse(dateStr.Replace("UTC", "").Trim(), out var createdAt);
|
||||||
|
|
||||||
|
return new TelegramMessage(id, senderName, createdAt, content);
|
||||||
|
}
|
||||||
|
catch { return null; }
|
||||||
|
}
|
||||||
|
|
||||||
|
private string CleanName(IElement node)
|
||||||
|
{
|
||||||
|
var clone = (IElement)node.Clone();
|
||||||
|
foreach (var span in clone.QuerySelectorAll("span")) span.Remove();
|
||||||
|
return clone.TextContent.Trim();
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,12 +1,11 @@
|
|||||||
using Knot.Shared.Kernel;
|
using Knot.Shared.Kernel;
|
||||||
using Knot.Shared.Kernel.Configuration;
|
using Knot.Shared.Kernel.Configuration;
|
||||||
using Microsoft.Extensions.Configuration;
|
|
||||||
using MediatR;
|
using MediatR;
|
||||||
using System.Collections.Generic;
|
using System.Collections.Generic;
|
||||||
using System.Threading;
|
using System.Threading;
|
||||||
using System.Threading.Tasks;
|
using System.Threading.Tasks;
|
||||||
|
|
||||||
namespace Host.Application.WebRtc.Queries;
|
namespace Knot.Modules.WebRtc.Application.WebRtc.Queries;
|
||||||
|
|
||||||
public record IceServerDto(string[] Urls, string? Username = null, string? Credential = null);
|
public record IceServerDto(string[] Urls, string? Username = null, string? Credential = null);
|
||||||
public record IceServersResultDto(List<IceServerDto> IceServers);
|
public record IceServersResultDto(List<IceServerDto> IceServers);
|
||||||
@@ -15,59 +14,52 @@ public record GetIceServersQuery : IQuery<IceServersResultDto>;
|
|||||||
|
|
||||||
internal sealed class GetIceServersQueryHandler : IQueryHandler<GetIceServersQuery, IceServersResultDto>
|
internal sealed class GetIceServersQueryHandler : IQueryHandler<GetIceServersQuery, IceServersResultDto>
|
||||||
{
|
{
|
||||||
private readonly IConfiguration _configuration;
|
private readonly ISettingsService _settingsService;
|
||||||
private readonly IWebRtcSettings _settingsService;
|
|
||||||
|
|
||||||
public GetIceServersQueryHandler(IConfiguration configuration, IWebRtcSettings settingsService)
|
public GetIceServersQueryHandler(ISettingsService settingsService)
|
||||||
{
|
{
|
||||||
_configuration = configuration;
|
|
||||||
_settingsService = settingsService;
|
_settingsService = settingsService;
|
||||||
}
|
}
|
||||||
|
|
||||||
public Task<Result<IceServersResultDto>> Handle(GetIceServersQuery request, CancellationToken cancellationToken)
|
public async Task<Result<IceServersResultDto>> Handle(GetIceServersQuery request, CancellationToken cancellationToken)
|
||||||
{
|
{
|
||||||
var settings = _settingsService.Current;
|
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
|
||||||
if (!settings.Enabled)
|
var webrtc = settings.WebRtc;
|
||||||
|
|
||||||
|
if (!webrtc.Enabled)
|
||||||
{
|
{
|
||||||
return Task.FromResult(Result.Failure<IceServersResultDto>(new Error(
|
return Result.Failure<IceServersResultDto>(new Error(
|
||||||
Knot.Shared.Kernel.Constants.Errors.DisabledByAdmin,
|
"WebRtc.Disabled",
|
||||||
"СервРСвЂВВР РЋР С“ отключен Р°РТвЂВВР В Р’В Р РЋР’ВВР В Р’В Р РЋРІР‚ВВР Р…Р СвЂВВстратороРСВВВ."
|
"WebRTC calls are disabled by the server administrator."));
|
||||||
)));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var turnUrl = !string.IsNullOrEmpty(settings.TurnHost)
|
|
||||||
? $"turn:{settings.TurnHost}:{settings.TurnPort}"
|
|
||||||
: _configuration["WebRtc:TurnUrl"];
|
|
||||||
|
|
||||||
var turnUsername = !string.IsNullOrEmpty(settings.TurnUser)
|
|
||||||
? settings.TurnUser
|
|
||||||
: _configuration["WebRtc:TurnUsername"];
|
|
||||||
|
|
||||||
var turnSecret = !string.IsNullOrEmpty(settings.TurnSecret)
|
|
||||||
? settings.TurnSecret
|
|
||||||
: _configuration["WebRtc:TurnPassword"];
|
|
||||||
|
|
||||||
var iceServers = new List<IceServerDto>();
|
var iceServers = new List<IceServerDto>();
|
||||||
|
|
||||||
if (!string.IsNullOrEmpty(turnUrl))
|
// Если хост TURN задан, формируем STUN и TURN записи
|
||||||
|
if (!string.IsNullOrEmpty(webrtc.TurnHost))
|
||||||
{
|
{
|
||||||
var stunUrl = turnUrl.Replace("turn:", "stun:");
|
var turnUrl = $"turn:{webrtc.TurnHost}:{webrtc.TurnPort}";
|
||||||
iceServers.Add(new IceServerDto(new[] { stunUrl }));
|
var stunUrl = $"stun:{webrtc.TurnHost}:{webrtc.TurnPort}";
|
||||||
|
|
||||||
if (!string.IsNullOrEmpty(turnUsername))
|
// Всегда добавляем STUN (публичный или свой)
|
||||||
|
iceServers.Add(new IceServerDto(new[] { stunUrl, "stun:stun.l.google.com:19302" }));
|
||||||
|
|
||||||
|
// Добавляем TURN (TCP и UDP варианты)
|
||||||
|
if (!string.IsNullOrEmpty(webrtc.TurnUser))
|
||||||
{
|
{
|
||||||
iceServers.Add(new IceServerDto(
|
iceServers.Add(new IceServerDto(
|
||||||
new[] { turnUrl, turnUrl + "?transport=tcp" },
|
new[] { turnUrl, $"{turnUrl}?transport=tcp" },
|
||||||
turnUsername,
|
webrtc.TurnUser,
|
||||||
!string.IsNullOrEmpty(turnSecret) ? turnSecret : turnUsername
|
!string.IsNullOrEmpty(webrtc.TurnSecret) ? webrtc.TurnSecret : webrtc.TurnUser
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
else
|
}
|
||||||
{
|
else
|
||||||
iceServers.Add(new IceServerDto(new[] { turnUrl, turnUrl + "?transport=tcp" }));
|
{
|
||||||
}
|
// Если своего сервера нет, используем публичный Google STUN как fallback
|
||||||
|
iceServers.Add(new IceServerDto(new[] { "stun:stun.l.google.com:19302" }));
|
||||||
}
|
}
|
||||||
|
|
||||||
return Task.FromResult(Result.Success(new IceServersResultDto(iceServers)));
|
return Result.Success(new IceServersResultDto(iceServers));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1,25 +1,26 @@
|
|||||||
using Carter;
|
using Carter;
|
||||||
using MediatR;
|
using MediatR;
|
||||||
using Knot.Shared.Kernel;
|
using Knot.Shared.Kernel;
|
||||||
|
using Knot.Modules.WebRtc.Application.WebRtc.Queries;
|
||||||
using Microsoft.AspNetCore.Builder;
|
using Microsoft.AspNetCore.Builder;
|
||||||
using Microsoft.AspNetCore.Http;
|
using Microsoft.AspNetCore.Http;
|
||||||
using Microsoft.AspNetCore.Routing;
|
using Microsoft.AspNetCore.Routing;
|
||||||
|
|
||||||
namespace Host.Endpoints;
|
namespace Knot.Modules.WebRtc.Presentation.Endpoints;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// РегРСвЂВВстрацРСвЂВВР РЋР РЏ РЎРЊР Р…Р ТвЂВВРїРѕРСвЂВВнтовWebRTC
|
/// Endpoints for WebRTC signaling and ICE configuration
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public sealed class WebRtcEndpoints : ICarterModule
|
public sealed class WebRtcEndpoints : ICarterModule
|
||||||
{
|
{
|
||||||
public void AddRoutes(IEndpointRouteBuilder app)
|
public void AddRoutes(IEndpointRouteBuilder app)
|
||||||
{
|
{
|
||||||
var group = app.MapGroup(Knot.Shared.Kernel.Constants.Routes.ApiWebRtc)
|
var group = app.MapGroup("api/webrtc")
|
||||||
.RequireAuthorization();
|
.RequireAuthorization();
|
||||||
|
|
||||||
group.MapGet("/ice-servers", async (ISender sender, CancellationToken ct) =>
|
group.MapGet("/ice-servers", async (ISender sender, CancellationToken ct) =>
|
||||||
{
|
{
|
||||||
var result = await sender.Send(new Host.Application.WebRtc.Queries.GetIceServersQuery(), ct);
|
var result = await sender.Send(new GetIceServersQuery(), ct);
|
||||||
return result.IsSuccess ? Results.Ok(result.Value) : Results.BadRequest(result.Error);
|
return result.IsSuccess ? Results.Ok(result.Value) : Results.BadRequest(result.Error);
|
||||||
})
|
})
|
||||||
.WithName("GetIceServers")
|
.WithName("GetIceServers")
|
||||||
|
|||||||
Reference in New Issue
Block a user