Структура, доп модули, федерация, документация

This commit is contained in:
Халимов Рустам
2026-03-27 00:55:01 +03:00
parent 7cb6ac61dd
commit 7ef73b414c
64 changed files with 3080 additions and 133 deletions

View File

@@ -0,0 +1,38 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
namespace Knot.Modules.Federation.Application.Abstractions;
/// <summary>
/// Шлюз для безопасной передачи данных между серверами Knot Messager.
/// </summary>
public interface IFederationGateway
{
/// <summary>
/// Отправляет пакет на конкретный домен.
/// </summary>
Task<Result> SendPacketAsync(FederationMessagePacket packet, string targetDomain, CancellationToken ct = default);
}
/// <summary>
/// Федеративный пакет сообщения с гибридным шифрованием.
/// </summary>
public record FederationMessagePacket(
string SenderDomain,
string RecipientDomain,
string EncryptedPayload, // AES-256
string EncryptedIV, // RSA-Encrypted for Recipient
string Signature, // RSA-Signed by Sender
FederationMetadata Metadata
);
public record FederationMetadata(
Guid ChatId,
Guid SenderId,
string SenderUsername,
string MessageType,
DateTime CreatedAt,
Guid UserId = default // ID пользователя-исполнителя (для реакций и т.д.)
);

View File

@@ -1,54 +1,66 @@
using Knot.Modules.Stories.Domain;
using Knot.Shared.Kernel.Constants;
using System.Security.Cryptography;
using MediatR;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Configuration;
using Knot.Shared.Kernel.Constants;
using System;
using System.Linq;
using System.Security.Cryptography;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
namespace Host.Application.Federation.Commands;
public record HandshakeRequest(string Domain, string PublicKey);
public record HandshakeResponse(string Domain, string PublicKey, string Status);
public record HandshakeRequest(string Domain, string PublicKey, RemoteCapabilities Capabilities);
public record HandshakeResponse(string Domain, string PublicKey, RemoteCapabilities Capabilities, string Status);
public record HandshakeFederationCommand(HandshakeRequest Request) : ICommand<HandshakeResponse>;
internal sealed class HandshakeFederationCommandHandler : ICommandHandler<HandshakeFederationCommand, HandshakeResponse>
{
private readonly ISettingsService _settings;
private readonly ISettingsService _settingsService;
public HandshakeFederationCommandHandler(ISettingsService settings)
public HandshakeFederationCommandHandler(ISettingsService settingsService)
{
_settings = settings;
_settingsService = settingsService;
}
public Task<Result<HandshakeResponse>> Handle(HandshakeFederationCommand request, CancellationToken cancellationToken)
public async Task<Result<HandshakeResponse>> Handle(HandshakeFederationCommand request, CancellationToken cancellationToken)
{
var conf = _settings.Current;
if (!conf.EnableConfederation)
return Task.FromResult(Result.Failure<HandshakeResponse>(new Error(Errors.DisabledByAdmin, "Federation is disabled")));
var allSettings = await _settingsService.GetSettingsAsync(cancellationToken);
var conf = allSettings.Federation;
if (!conf.Enabled)
return Result.Failure<HandshakeResponse>(new Error(Errors.DisabledByAdmin, "Federation is disabled or keys not generated."));
if (string.IsNullOrEmpty(request.Request.Domain) || string.IsNullOrEmpty(request.Request.PublicKey))
return Task.FromResult(Result.Failure<HandshakeResponse>(DomainErrors.RequestInvalid));
return Result.Failure<HandshakeResponse>(new Error("Request.Invalid", "Domain and PublicKey required."));
var allowedList = conf.AllowedDomains?.Select(d => d.Trim().ToLower()).ToList() ?? new System.Collections.Generic.List<string>();
if (!allowedList.Contains(request.Request.Domain.ToLowerInvariant()))
return Task.FromResult(Result.Failure<HandshakeResponse>(StoryErrors.Unauthorized));
// Проверка разрешенности домена (белый список)
var allowedList = conf.AllowedDomains ?? new List<FederationDomainConfig>();
var domainConfig = allowedList.FirstOrDefault(d => d.IsEnabled && d.Domain.Equals(request.Request.Domain, StringComparison.OrdinalIgnoreCase));
if (domainConfig == null)
return Result.Failure<HandshakeResponse>(new Error("Federation.DomainNotAllowed", "Your domain is not in our allowlist."));
using var rsa = RSA.Create(2048);
var selfPublicKey = Convert.ToBase64String(rsa.ExportRSAPublicKey());
// Сохраняем полученные данные о партнере (Capabilities и Ключ) в кэш настроек
domainConfig.PublicKey = request.Request.PublicKey;
domainConfig.Capabilities = request.Request.Capabilities;
await _settingsService.UpdateSettingsAsync(allSettings, cancellationToken);
// Наш ответ с нашими возможностями
var ourCapabilities = new RemoteCapabilities
{
AllowMedia = allSettings.Messages.AllowMedia,
AllowPolls = allSettings.Messages.AllowPolls,
AllowVoiceMessages = allSettings.Messages.AllowVoiceMessages,
AllowVideoCalls = allSettings.WebRtc.EnableVideoCalls,
AllowScreenSharing = allSettings.WebRtc.EnableScreenSharing
};
var response = new HandshakeResponse(
Environment.GetEnvironmentVariable("DOMAIN") ?? "knot.local",
selfPublicKey,
allSettings.System.DomainUrl,
conf.PublicKey ?? string.Empty,
ourCapabilities,
"Accepted"
);
return Task.FromResult(Result.Success(response));
return Result.Success(response);
}
}

View File

@@ -0,0 +1,206 @@
using System;
using System.Linq;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel;
using Knot.Modules.Auth.Domain;
using Knot.Modules.Messaging.Application.Abstractions;
using Knot.Modules.Messaging.Domain;
using Knot.Shared.Kernel.Configuration;
using Knot.Shared.Kernel.Constants;
namespace Host.Application.Federation.Commands;
public record InboundFederationCommand(FederationMessagePacket Packet) : ICommand;
internal sealed class InboundFederationCommandHandler : ICommandHandler<InboundFederationCommand>
{
private readonly ISettingsService _settingsService;
private readonly IMessageRepository _messageRepository;
private readonly IUserRepository _userRepository;
private readonly IMessageReactionRepository _reactionRepository;
private readonly IMessageNotifier _notifier;
private readonly IMediator _mediator;
public InboundFederationCommandHandler(
ISettingsService settingsService,
IMessageRepository messageRepository,
IUserRepository userRepository,
IMessageReactionRepository reactionRepository,
IMessageNotifier notifier,
IMediator mediator)
{
_settingsService = settingsService;
_messageRepository = messageRepository;
_userRepository = userRepository;
_reactionRepository = reactionRepository;
_notifier = notifier;
_mediator = mediator;
}
public async Task<Result> Handle(InboundFederationCommand request, CancellationToken cancellationToken)
{
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
if (!settings.Federation.Enabled)
return Result.Failure(new Error(Errors.DisabledByAdmin, "Federation disabled."));
var packet = request.Packet;
// 1. Проверяем подпись отправителя
var senderConfig = settings.Federation.AllowedDomains.FirstOrDefault(d => d.IsEnabled && d.Domain.Equals(packet.SenderDomain, StringComparison.OrdinalIgnoreCase));
if (senderConfig == null || string.IsNullOrEmpty(senderConfig.PublicKey))
return Result.Failure(new Error("Federation.SenderNotAllowed", "Sender domain not in allowlist."));
using (var rsaVerify = RSA.Create())
{
rsaVerify.ImportRSAPublicKey(Convert.FromBase64String(senderConfig.PublicKey), out _);
var dataToVerify = Encoding.UTF8.GetBytes(packet.EncryptedPayload + packet.SenderDomain + packet.RecipientDomain);
if (!rsaVerify.VerifyData(dataToVerify, Convert.FromBase64String(packet.Signature), HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1))
{
return Result.Failure(new Error("Federation.InvalidSignature", "RSA Signature verification failed."));
}
}
// 2. Расшифровываем IV своим Private Key
byte[] aesKey;
byte[] aesIv;
using (var rsaDecrypt = RSA.Create())
{
rsaDecrypt.ImportPkcs8PrivateKey(Convert.FromBase64String(settings.Federation.PrivateKey!), out _);
var decryptedKeys = rsaDecrypt.Decrypt(Convert.FromBase64String(packet.EncryptedIV), RSAEncryptionPadding.Pkcs1);
// AES-256 Key (32 bytes) + IV (16 bytes)
aesKey = new byte[32];
aesIv = new byte[16];
Buffer.BlockCopy(decryptedKeys, 0, aesKey, 0, 32);
Buffer.BlockCopy(decryptedKeys, 32, aesIv, 0, 16);
}
// 3. Расшифровываем само сообщение (AES-256)
string plainText;
using (var aes = Aes.Create())
{
aes.Key = aesKey;
aes.IV = aesIv;
using (var decryptor = aes.CreateDecryptor())
{
var cipherBytes = Convert.FromBase64String(packet.EncryptedPayload);
var plainBytes = decryptor.TransformFinalBlock(cipherBytes, 0, cipherBytes.Length);
plainText = Encoding.UTF8.GetString(plainBytes);
}
}
// 4. Логика обработки типов сообщений
if (packet.Metadata.MessageType == "sync_capabilities")
{
var capabilities = JsonSerializer.Deserialize<RemoteCapabilities>(plainText);
if (capabilities != null && senderConfig != null)
{
senderConfig.Capabilities = capabilities;
await _settingsService.UpdateSettingsAsync(settings, cancellationToken);
}
return Result.Success();
}
if (packet.Metadata.MessageType == "presence_update")
{
var statusData = JsonSerializer.Deserialize<JsonElement>(plainText);
var isOnline = statusData.GetProperty("IsOnline").GetBoolean();
var user = await _userRepository.GetByIdAsync(packet.Metadata.SenderId, cancellationToken);
if (user != null && user.IsExternal)
{
user.UpdateStatus(isOnline);
_userRepository.Update(user);
// Уведомляем локальных пользователей через SignalR
await _notifier.NotifyNewMessageAsync(Guid.Empty, new { type = "presence_update", userId = user.Id, isOnline = user.IsOnline }, cancellationToken);
}
return Result.Success();
}
if (packet.Metadata.MessageType == "message_edited")
{
var messageToEdit = await _messageRepository.GetByIdAsync(packet.Metadata.SenderId, cancellationToken); // В метаданных MessageId
if (messageToEdit != null)
{
messageToEdit.Edit(plainText);
await _messageRepository.UpdateAsync(messageToEdit, cancellationToken);
await _notifier.NotifyNewMessageAsync(messageToEdit.ChatId, new { type = "message_edited", messageId = messageToEdit.Id, content = plainText }, cancellationToken);
}
return Result.Success();
}
if (packet.Metadata.MessageType == "message_deleted")
{
var messageToDelete = await _messageRepository.GetByIdAsync(packet.Metadata.SenderId, cancellationToken);
if (messageToDelete != null)
{
messageToDelete.Delete();
await _messageRepository.UpdateAsync(messageToDelete, cancellationToken);
await _notifier.NotifyNewMessageAsync(messageToDelete.ChatId, new { type = "message_deleted", messageId = messageToDelete.Id }, cancellationToken);
}
return Result.Success();
}
if (packet.Metadata.MessageType == "reaction_added" || packet.Metadata.MessageType == "reaction_removed")
{
var messageForReaction = await _messageRepository.GetByIdAsync(packet.Metadata.SenderId, cancellationToken);
if (messageForReaction != null)
{
if (packet.Metadata.MessageType == "reaction_added")
{
var reaction = new MessageReaction(messageForReaction.Id, packet.Metadata.UserId, plainText);
await _reactionRepository.AddAsync(reaction, cancellationToken);
}
else
{
await _reactionRepository.RemoveAsync(messageForReaction.Id, packet.Metadata.UserId, plainText, cancellationToken);
}
await _notifier.NotifyNewMessageAsync(messageForReaction.ChatId, new {
type = packet.Metadata.MessageType,
messageId = messageForReaction.Id,
userId = packet.Metadata.UserId,
emoji = plainText
}, cancellationToken);
}
return Result.Success();
}
if (packet.Metadata.MessageType == "rtc_signal")
{
// Здесь проброс WebRTC сигнала (Offer/Answer/ICE) конечному пользователю через SignalR
await _notifier.NotifyNewMessageAsync(packet.Metadata.ChatId, new { type = "rtc_signal", payload = plainText, senderId = packet.Metadata.SenderId }, cancellationToken);
return Result.Success();
}
if (packet.Metadata.MessageType == "poll" && !settings.Messages.AllowPolls)
return Result.Failure(new Error("Federation.PollsDisabled", "We do not accept polls."));
// 5. Сохранение в базу сообщений (MongoDB)
// В реальном приложении здесь будет маппинг на TextMessage, MediaMessage и т.д.
var message = new TextMessage(
Guid.NewGuid(),
packet.Metadata.ChatId,
packet.Metadata.SenderId,
plainText,
null, // replyToId
null, // quote
null, // forwardedFromId
packet.Metadata.CreatedAt,
false); // isImported
_messageRepository.Add(message);
// 6. Уведомление пользователя через SignalR
await _notifier.NotifyNewMessageAsync(packet.Metadata.ChatId, message, cancellationToken);
return Result.Success();
}
}

View File

@@ -0,0 +1,40 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Shared.Kernel;
using Knot.Modules.Auth.Domain;
using Knot.Shared.Kernel.Configuration;
using System.Linq;
namespace Host.Application.Federation.Commands;
public record ResolveUserResponse(Guid Id, string Username, string DisplayName, string? Avatar);
public record ResolveUserCommand(string Username) : ICommand<ResolveUserResponse>;
internal sealed class ResolveUserCommandHandler : ICommandHandler<ResolveUserCommand, ResolveUserResponse>
{
private readonly IUserRepository _userRepository;
public ResolveUserCommandHandler(IUserRepository userRepository)
{
_userRepository = userRepository;
}
public async Task<Result<ResolveUserResponse>> Handle(ResolveUserCommand request, CancellationToken cancellationToken)
{
var user = await _userRepository.GetByUsernameAsync(request.Username, cancellationToken);
if (user == null || user.IsExternal)
{
return Result.Failure<ResolveUserResponse>(new Error("Federation.UserNotFound", "User not found on this server."));
}
return Result.Success(new ResolveUserResponse(
user.Id,
user.Username,
user.DisplayName,
user.Avatar
));
}
}

View File

@@ -0,0 +1,97 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Modules.Messaging.Domain;
using Knot.Modules.Federation.Application.Federation.Services;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel.Configuration;
using Knot.Modules.Auth.Domain;
namespace Knot.Modules.Federation.Application.Federation.Events;
/// <summary>
/// Универсальный обработчик действий с сообщениями (Edit, Delete, Reactions) для Федерации.
/// </summary>
public sealed class FederatedMessageActionsHandler :
INotificationHandler<MessageEditedDomainEvent>,
INotificationHandler<MessageDeletedDomainEvent>,
INotificationHandler<MessageReactionAddedDomainEvent>,
INotificationHandler<MessageReactionRemovedDomainEvent>
{
private readonly IChatRepository _chatRepository;
private readonly IUserRepository _userRepository;
private readonly ISettingsService _settingsService;
private readonly FederationPacketService _packetService;
private readonly IFederationGateway _gateway;
public FederatedMessageActionsHandler(
IChatRepository chatRepository,
IUserRepository userRepository,
ISettingsService settingsService,
FederationPacketService packetService,
IFederationGateway gateway)
{
_chatRepository = chatRepository;
_userRepository = userRepository;
_settingsService = settingsService;
_packetService = packetService;
_gateway = gateway;
}
public async Task Handle(MessageEditedDomainEvent notification, CancellationToken cancellationToken)
=> await ProcessAction(notification.ChatId, notification.MessageId, "message_edited", notification.NewContent, cancellationToken);
public async Task Handle(MessageDeletedDomainEvent notification, CancellationToken cancellationToken)
=> await ProcessAction(notification.ChatId, notification.MessageId, "message_deleted", string.Empty, cancellationToken);
public async Task Handle(MessageReactionAddedDomainEvent notification, CancellationToken cancellationToken)
=> await ProcessAction(notification.ChatId, notification.MessageId, "reaction_added", notification.Emoji, cancellationToken, notification.UserId);
public async Task Handle(MessageReactionRemovedDomainEvent notification, CancellationToken cancellationToken)
=> await ProcessAction(notification.ChatId, notification.MessageId, "reaction_removed", notification.Emoji, cancellationToken, notification.UserId);
private async Task ProcessAction(Guid chatId, Guid messageId, string actionType, string payload, CancellationToken cancellationToken, Guid actorId = default)
{
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
if (!settings.Federation.Enabled) return;
var chat = await _chatRepository.GetByIdAsync(chatId, cancellationToken);
if (chat == null) return;
var externalDomains = await GetExternalDomains(chat, cancellationToken);
if (!externalDomains.Any()) return;
var metadata = new FederationMetadata(
chatId,
messageId, // MessageId в SenderId для операций
"system",
actionType,
DateTime.UtcNow,
actorId
);
foreach (var domain in externalDomains)
{
var packetResult = await _packetService.PreparePacketAsync(payload, metadata, domain, cancellationToken);
if (packetResult.IsSuccess)
{
await _gateway.SendPacketAsync(packetResult.Value, domain, cancellationToken);
}
}
}
private async Task<List<string>> GetExternalDomains(Chat chat, CancellationToken cancellationToken)
{
var memberIds = chat.Members.Select(m => m.UserId).ToList();
var members = await _userRepository.GetByIdsAsync(memberIds, cancellationToken);
return members
.Where(u => u.IsExternal && !string.IsNullOrEmpty(u.Domain))
.Select(u => u.Domain!)
.Distinct()
.ToList();
}
}

View File

@@ -0,0 +1,90 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Modules.Messaging.Domain;
using Knot.Modules.Federation.Application.Federation.Services;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel.Configuration;
using Knot.Modules.Auth.Domain;
namespace Knot.Modules.Federation.Application.Federation.Events;
/// <summary>
/// Обработчик доменного события отправки сообщения.
/// Если в чате есть внешние участники — инициирует федеративную рассылку.
/// </summary>
public sealed class MessageSentDomainEventHandler : INotificationHandler<MessageSentDomainEvent>
{
private readonly IChatRepository _chatRepository;
private readonly IUserRepository _userRepository;
private readonly ISettingsService _settingsService;
private readonly FederationPacketService _packetService;
private readonly IFederationGateway _gateway;
public MessageSentDomainEventHandler(
IChatRepository chatRepository,
IUserRepository userRepository,
ISettingsService settingsService,
FederationPacketService packetService,
IFederationGateway gateway)
{
_chatRepository = chatRepository;
_userRepository = userRepository;
_settingsService = settingsService;
_packetService = packetService;
_gateway = gateway;
}
public async Task Handle(MessageSentDomainEvent notification, CancellationToken cancellationToken)
{
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
if (!settings.Federation.Enabled) return;
// 1. Загружаем чат и проверяем наличие внешних участников
var chat = await _chatRepository.GetByIdAsync(notification.ChatId, cancellationToken);
if (chat == null) return;
var memberIds = chat.Members.Select(m => m.UserId).ToList();
var members = await _userRepository.GetByIdsAsync(memberIds, cancellationToken);
// Находим уникальные домены внешних участников
var externalDomains = members
.Where(u => u.IsExternal && !string.IsNullOrEmpty(u.Domain))
.Select(u => u.Domain!)
.Distinct()
.ToList();
if (!externalDomains.Any()) return;
// 2. Получаем имя отправителя для метаданных
var sender = await _userRepository.GetByIdAsync(notification.SenderId, cancellationToken);
var senderUsername = sender?.Username ?? "unknown";
var metadata = new FederationMetadata(
notification.ChatId,
notification.SenderId,
senderUsername,
"text", // Для начала поддерживаем только текст через ивент
DateTime.UtcNow
);
// 3. Рассылка по доменам (Fan-out)
foreach (var domain in externalDomains)
{
var packetResult = await _packetService.PreparePacketAsync(
notification.Content ?? string.Empty,
metadata,
domain,
cancellationToken);
if (packetResult.IsSuccess)
{
// Отправляем асинхронно
await _gateway.SendPacketAsync(packetResult.Value, domain, cancellationToken);
}
}
}
}

View File

@@ -0,0 +1,71 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Modules.Admin.Domain.Events;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Modules.Federation.Application.Federation.Services;
using Knot.Shared.Kernel.Configuration;
namespace Knot.Modules.Federation.Application.Federation.Events;
/// <summary>
/// Обработчик события обновления настроек.
/// Уведомляет все удаленные серверы о смене способностей (Capabilities).
/// </summary>
public sealed class SystemSettingsUpdatedDomainEventHandler : INotificationHandler<SystemSettingsUpdatedDomainEvent>
{
private readonly ISettingsService _settingsService;
private readonly IFederationGateway _gateway;
private readonly FederationPacketService _packetService;
public SystemSettingsUpdatedDomainEventHandler(
ISettingsService settingsService,
IFederationGateway gateway,
FederationPacketService packetService)
{
_settingsService = settingsService;
_gateway = gateway;
_packetService = packetService;
}
public async Task Handle(SystemSettingsUpdatedDomainEvent notification, CancellationToken cancellationToken)
{
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
if (!settings.Federation.Enabled) return;
// Наши новые способности
var ourCapabilities = new RemoteCapabilities
{
AllowMedia = settings.Messages.AllowMedia,
AllowPolls = settings.Messages.AllowPolls,
AllowVoiceMessages = settings.Messages.AllowVoiceMessages,
AllowVideoCalls = settings.WebRtc.EnableVideoCalls,
AllowScreenSharing = settings.WebRtc.EnableScreenSharing
};
// Собираем метаданные для синхронизации
var metadata = new FederationMetadata(
Guid.Empty, // Системный чат
Guid.Empty, // Системный пользователь
"SYSTEM",
"sync_capabilities",
DateTime.UtcNow
);
// Рассылка по всем доменам из белого списка (Fan-out)
foreach (var domainConfig in settings.Federation.AllowedDomains.Where(d => d.IsEnabled))
{
// Упаковываем Capabilities в JSON для передачи в payload
var payload = System.Text.Json.JsonSerializer.Serialize(ourCapabilities);
var packetResult = await _packetService.PreparePacketAsync(payload, metadata, domainConfig.Domain, cancellationToken);
if (packetResult.IsSuccess)
{
await _gateway.SendPacketAsync(packetResult.Value, domainConfig.Domain, cancellationToken);
}
}
}
}

View File

@@ -0,0 +1,86 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using Knot.Modules.Auth.Domain;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Modules.Federation.Application.Federation.Services;
using Knot.Shared.Kernel.Configuration;
namespace Knot.Modules.Federation.Application.Federation.Events;
/// <summary>
/// Обработчик изменения статуса пользователя.
/// Рассылает новый статус всем серверам, где у пользователя есть чаты.
/// </summary>
public sealed class UserStatusChangedDomainEventHandler : INotificationHandler<UserStatusChangedDomainEvent>
{
private readonly IChatRepository _chatRepository;
private readonly IUserRepository _userRepository;
private readonly ISettingsService _settingsService;
private readonly FederationPacketService _packetService;
private readonly IFederationGateway _gateway;
public UserStatusChangedDomainEventHandler(
IChatRepository chatRepository,
IUserRepository userRepository,
ISettingsService settingsService,
FederationPacketService packetService,
IFederationGateway gateway)
{
_chatRepository = chatRepository;
_userRepository = userRepository;
_settingsService = settingsService;
_packetService = packetService;
_gateway = gateway;
}
public async Task Handle(UserStatusChangedDomainEvent notification, CancellationToken cancellationToken)
{
var settings = await _settingsService.GetSettingsAsync(cancellationToken);
if (!settings.Federation.Enabled) return;
// 1. Находим все чаты пользователя
var userChats = await _chatRepository.GetUserChatsAsync(notification.UserId, cancellationToken);
if (!userChats.Any()) return;
// 2. Определяем уникальные внешние домены участников этих чатов
var allMemberIds = userChats.SelectMany(c => c.Members.Select(m => m.UserId)).Distinct().ToList();
var members = await _userRepository.GetByIdsAsync(allMemberIds, cancellationToken);
var externalDomains = members
.Where(u => u.IsExternal && !string.IsNullOrEmpty(u.Domain))
.Select(u => u.Domain!)
.Distinct()
.ToList();
if (!externalDomains.Any()) return;
// 3. Формируем пакет статуса
var statusData = new {
notification.IsOnline,
notification.LastSeen
};
var payload = System.Text.Json.JsonSerializer.Serialize(statusData);
var metadata = new FederationMetadata(
Guid.Empty,
notification.UserId,
"system", // Username отправителя здесь не важен, важен SenderId
"presence_update",
DateTime.UtcNow
);
// 4. Рассылка по доменам
foreach (var domain in externalDomains)
{
var packetResult = await _packetService.PreparePacketAsync(payload, metadata, domain, cancellationToken);
if (packetResult.IsSuccess)
{
await _gateway.SendPacketAsync(packetResult.Value, domain, cancellationToken);
}
}
}
}

View File

@@ -0,0 +1,107 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Configuration;
using Knot.Shared.Kernel.Security;
namespace Knot.Modules.Federation.Application.Federation.Services;
/// <summary>
/// Сервис для подготовки и рассылки зашифрованных федератовных пакетов.
/// </summary>
public sealed class FederationPacketService
{
private readonly ISettingsService _settingsService;
private readonly IFederationGateway _gateway;
private readonly IEncryptionService _encryption;
public FederationPacketService(
ISettingsService settingsService,
IFederationGateway gateway,
IEncryptionService encryption)
{
_settingsService = settingsService;
_gateway = gateway;
_encryption = encryption;
}
/// <summary>
/// Собирает зашифрованный пакет для целевого домена.
/// </summary>
public async Task<Result<FederationMessagePacket>> PreparePacketAsync(
string messageBody,
FederationMetadata metadata,
string targetDomain,
CancellationToken ct = default)
{
var settings = await _settingsService.GetSettingsAsync(ct);
var ourDomain = settings.System.DomainUrl;
var ourPrivateKeyBase64 = settings.Federation.PrivateKey;
// Находим ключ получателя в белом списке
var targetConfig = settings.Federation.AllowedDomains.FirstOrDefault(d => d.Domain.Equals(targetDomain, StringComparison.OrdinalIgnoreCase));
if (targetConfig == null || string.IsNullOrEmpty(targetConfig.PublicKey))
{
return Result.Failure<FederationMessagePacket>(new Error("Federation.KeyNotFound", "Public key for target domain not found."));
}
// 1. AES Шифрование тела
// Используем IEncryptionService для получения потока или метода шифрования (в реальности здесь будет AES-256)
// Для демонстрации представим, что мы получили IV и EncryptedPayload
using var aes = Aes.Create();
aes.KeySize = 256;
aes.GenerateKey();
aes.GenerateIV();
var iv = aes.IV;
var key = aes.Key; // В гибридной схеме мы передаем IV + Key (или только IV при общем ключе)
byte[] encryptedBytes;
using (var encryptor = aes.CreateEncryptor())
{
var plainBytes = Encoding.UTF8.GetBytes(messageBody);
encryptedBytes = encryptor.TransformFinalBlock(plainBytes, 0, plainBytes.Length);
}
// 2. RSA упаковка ключей на Open Public Key получателя
byte[] encryptedKeys;
using (var rsa = RSA.Create())
{
rsa.ImportRSAPublicKey(Convert.FromBase64String(targetConfig.PublicKey), out _);
// Шифруем Key + IV (48 байт)
var combinedKeys = new byte[key.Length + iv.Length];
Buffer.BlockCopy(key, 0, combinedKeys, 0, key.Length);
Buffer.BlockCopy(iv, 0, combinedKeys, key.Length, iv.Length);
encryptedKeys = rsa.Encrypt(combinedKeys, RSAEncryptionPadding.Pkcs1);
}
// 3. Подпись всего пакета нашим Private Key
string signature;
using (var rsaSign = RSA.Create())
{
rsaSign.ImportPkcs8PrivateKey(Convert.FromBase64String(ourPrivateKeyBase64!), out _);
var dataToSign = Encoding.UTF8.GetBytes(Convert.ToBase64String(encryptedBytes) + ourDomain + targetDomain);
var sigBytes = rsaSign.SignData(dataToSign, HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1);
signature = Convert.ToBase64String(sigBytes);
}
var packet = new FederationMessagePacket(
ourDomain,
targetDomain,
Convert.ToBase64String(encryptedBytes),
Convert.ToBase64String(encryptedKeys),
signature,
metadata
);
return Result.Success(packet);
}
}

View File

@@ -0,0 +1,58 @@
using System;
using System.Net.Http;
using System.Net.Http.Json;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Configuration;
using Knot.Shared.Kernel.Security;
namespace Knot.Modules.Federation.Infrastructure.Services;
/// <summary>
/// Реализация шлюза федерации с гибридным шифрованием (AES-RSA).
/// </summary>
public sealed class FederationGateway : IFederationGateway
{
private readonly ISettingsService _settingsService;
private readonly HttpClient _httpClient;
public FederationGateway(ISettingsService settingsService, HttpClient httpClient)
{
_settingsService = settingsService;
_httpClient = httpClient;
}
public async Task<Result> SendPacketAsync(FederationMessagePacket packet, string targetDomain, CancellationToken ct = default)
{
try
{
var url = $"{targetDomain.TrimEnd('/')}/api/federation/v1/inbound";
// 1. Формируем тело запроса
var content = JsonContent.Create(packet);
// 2. Добавляем подпись в заголовки (как доп. уровень верификации)
_httpClient.DefaultRequestHeaders.Remove("X-Knot-Signature");
_httpClient.DefaultRequestHeaders.Add("X-Knot-Signature", packet.Signature);
var response = await _httpClient.PostAsync(url, content, ct);
if (response.IsSuccessStatusCode)
{
return Result.Success();
}
var errorMsg = await response.Content.ReadAsStringAsync(ct);
return Result.Failure(new Error("Federation.DeliveryFailed", $"Target server returned: {response.StatusCode}. {errorMsg}"));
}
catch (Exception ex)
{
return Result.Failure(new Error("Federation.NetworkError", ex.Message));
}
}
}

View File

@@ -6,11 +6,15 @@ using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Mvc;
using Microsoft.AspNetCore.Routing;
using System;
using Host.Application.Federation.Commands;
using Knot.Modules.Federation.Application.Abstractions;
using Knot.Shared.Kernel.Storage;
using Knot.Shared.Kernel.Configuration;
namespace Host.Endpoints;
/// <summary>
/// Регистрация эндпоинтов федерации.
/// Эндпоинты федерации для межсерверного взаимодействия.
/// </summary>
public sealed class FederationEndpoints : ICarterModule
{
@@ -18,24 +22,45 @@ public sealed class FederationEndpoints : ICarterModule
{
var group = app.MapGroup(Routes.ApiFederation);
group.MapPost("/handshake", async ([FromBody] Host.Application.Federation.Commands.HandshakeRequest request, ISender sender, CancellationToken ct) =>
group.MapPost("/handshake", async ([FromBody] HandshakeRequest request, ISender sender, CancellationToken ct) =>
{
var result = await sender.Send(new Host.Application.Federation.Commands.HandshakeFederationCommand(request), ct);
if (result.IsSuccess)
{
return Results.Ok(result.Value);
}
var result = await sender.Send(new HandshakeFederationCommand(request), ct);
return result.IsSuccess ? Results.Ok(result.Value) : Results.BadRequest(new { error = result.Error.Description });
});
if (result.Error.Code == "Unauthorized")
{
return Results.Forbid();
}
if (result.Error.Code == Knot.Shared.Kernel.Constants.Errors.DisabledByAdmin)
{
return Results.StatusCode(503);
}
group.MapPost("/inbound", async ([FromBody] FederationMessagePacket packet, ISender sender, CancellationToken ct) =>
{
var result = await sender.Send(new InboundFederationCommand(packet), ct);
return result.IsSuccess ? Results.Accepted() : Results.BadRequest(new { error = result.Error.Description });
});
group.MapGet("/resolve/{username}", async (string username, ISender sender, CancellationToken ct) =>
{
var result = await sender.Send(new ResolveUserCommand(username), ct);
return result.IsSuccess ? Results.Ok(result.Value) : Results.NotFound(new { error = result.Error.Description });
});
group.MapGet("/proxy/{id}", async (string id, [FromHeader(Name = "X-Knot-Signature")] string signature, [FromHeader(Name = "X-Knot-Domain")] string senderDomain, IFileStorageService storage, ISettingsService settingsService, CancellationToken ct) =>
{
// 1. Валидация подписи (упрощенно)
var settings = await settingsService.GetSettingsAsync(ct);
var senderConfig = settings.Federation.AllowedDomains.FirstOrDefault(d => d.Domain.Equals(senderDomain, StringComparison.OrdinalIgnoreCase));
return Results.BadRequest(new { error = result.Error.Description });
if (senderConfig == null || string.IsNullOrEmpty(senderConfig.PublicKey))
return Results.Forbid();
// В реальном коде здесь проверка RSA подписи параметров запроса (URL + Domain)
// 2. Стриминг файла первоисточника
try
{
var (stream, contentType, fileName) = await storage.DownloadFileAsync(id);
return Results.File(stream, contentType, fileName, enableRangeProcessing: true);
}
catch
{
return Results.NotFound();
}
});
}
}