40 Commits
Author SHA1 Message Date
max f4f61071ed внедрил пресеты статусов и модульную архитектуру 2026-04-18 00:43:46 +03:00
max e8c9a55fe9 безопасная обертка IsUserInvisibleAsync 2026-04-17 23:03:15 +03:00
max bccbf12a11 фикс <был недавно> 2026-04-17 22:24:09 +03:00
Халимов Рустам c8b0384392 Референс на контракты профиля 2026-04-17 21:07:09 +03:00
max 5245e8b7ae скрытие статуса online через сокеты и поле isInvisible в профиле 2026-04-17 20:47:28 +03:00
max b346372555 внедрена защищенная смена пароля (бэк + фронт) 2026-04-17 19:52:57 +03:00
max d2e6eb578a внедрил систему статусов (BIO) 2026-04-17 18:33:46 +03:00
max ae710c3d43 поправка групп 2026-04-17 17:56:30 +03:00
Халимов Рустам 600e43eec5 Merge remote-tracking branch 'origin/bugfix_web' into bugfix_web 2026-04-17 00:47:47 +03:00
Халимов Рустам 7b251686b2 Исправление подгрузки данных профиля 2026-04-17 00:43:07 +03:00
max 00682b2977 e# speciallyse enter a commit message to explain why this merge is
necessary,
if it merges an updated upstream into a topic branch.
2026-04-17 00:42:48 +03:00
max e05572fc3f Исправил неработющую функцию добавления пользователя через чат в группу 2026-04-17 00:42:03 +03:00
Халимов Рустам c264b7df27 исправил ошибки компиляции TypeScript 2026-04-17 00:37:37 +03:00
Халимов Рустам 56f75ae32b Получение данных профиля 2026-04-17 00:33:45 +03:00
Халимов Рустам ca9cf27716 О себе, редактирование в профиле 2026-04-17 00:26:38 +03:00
Халимов Рустам 7225e3272e Убрал аватары и имена внутри чата личных чатов 2026-04-16 23:54:19 +03:00
Халимов Рустам 86ae06beb6 Отступы в баблах 2026-04-16 23:48:36 +03:00
Халимов Рустам 454f70f716 Вставка и перетаскивание в поле ввода 2026-04-16 23:42:54 +03:00
Халимов Рустам e700609d30 Дубликат печати, разметка 2026-04-16 23:36:12 +03:00
Халимов Рустам 88f39aa51f Change for env 2026-04-16 22:42:26 +03:00
Халимов Рустам c3dbbaa7b8 Rename all 2026-04-16 22:30:37 +03:00
Халимов Рустам 2eb4f48ca0 Test create 2026-04-16 22:28:41 +03:00
Халимов Рустам 0c1adaab6c Rename web 2026-04-16 22:26:04 +03:00
Халимов Рустам a52726d0e6 Only web 2026-04-16 22:24:41 +03:00
Халимов Рустам d812a7a40c Change services name 2026-04-16 22:21:58 +03:00
Халимов Рустам c8b4fed25a Change ports 2026-04-16 22:18:04 +03:00
Халимов Рустам 8399d32490 Редактирование и плеер 2026-04-08 15:23:10 +03:00
Халимов Рустам 3905094ff4 Правка миграций 2026-04-07 21:22:38 +03:00
Халимов Рустам 9e8625aea1 Миграции 2026-04-07 21:20:06 +03:00
Халимов Рустам 5b905c94da Компоуз 2026-04-07 21:12:34 +03:00
Халимов Рустам 0eecc01374 Сборка 2026-04-07 21:05:23 +03:00
Халимов Рустам c37e1723d4 Фикс .env 2026-04-07 17:29:22 +03:00
Халимов Рустам 02043d4d97 Минимальная адаптация под мобилку 2026-04-07 12:29:41 +03:00
Халимов Рустам 852efa090e Опросы 2026-04-07 11:11:35 +03:00
Халимов Рустам c45f4db61c Заготовка опросов 2026-04-07 01:01:29 +03:00
Халимов Рустам 02a85fc587 Прокрутка 2026-04-06 23:42:46 +03:00
Халимов Рустам e09860700c Удален мусор 2026-04-06 23:36:42 +03:00
Халимов Рустам 32c9bc43cf Починка импорта, плеер для аудио 2026-04-06 23:35:45 +03:00
Халимов Рустам d96e4ec7d4 Правка импорта 2026-04-06 22:35:54 +03:00
Халимов Рустам 1558b20470 Импорт 2026-04-06 22:20:13 +03:00
174 changed files with 5348 additions and 5296 deletions
@@ -1,13 +1,13 @@
using FluentAssertions;
using Knot.Modules.Profiles.Application.Profiles.UpdateProfile;
using Knot.Modules.Profiles.Domain;
using Knot.Shared.Kernel;
using NSubstitute;
using Xunit;
using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Profiles.Application.Abstractions;
using Knot.Contracts.Profiles.Domain;
using Knot.Contracts.Profiles.Application.DTOs;
namespace Knot.Modules.Profiles.UnitTests;
@@ -25,14 +25,11 @@ public class UpdateProfileCommandHandlerTests
[Fact]
public async Task Handle_ShouldReturnError_WhenProfileNotFound()
{
// Arrange
var command = new UpdateProfileCommand(Guid.NewGuid(), "FirstName", "Bio", null);
_profileRepository.GetByIdAsync(command.UserId, Arg.Any<CancellationToken>()).Returns((ProfileDocument?)null);
var command = new UpdateProfileCommand(Guid.NewGuid(), "FirstName", "Bio", null, null);
_profileRepository.GetAsync(command.UserId, Arg.Any<CancellationToken>()).Returns((UserProfileDto?)null);
// Act
var result = await _handler.Handle(command, CancellationToken.None);
// Assert
result.IsFailure.Should().BeTrue();
}
}
@@ -8,4 +8,7 @@ public static class AuthErrors
public static Error IdentityRegistrationDisabled => new("Auth.RegistrationDisabled", "Registration is disabled");
public static Error IdentityUsernameNotUnique => new("Auth.UsernameNotUnique", "Username is already taken");
public static Error UserNotFound => new("Auth.UserNotFound", "User not found");
public static Error PasswordConfirmationMismatch => new("Auth.PasswordConfirmationMismatch", "Passwords do not match");
public static Error PasswordTooShort => new("Auth.PasswordTooShort", "Password must be at least 8 characters");
public static Error OldPasswordInvalid => new("Auth.OldPasswordInvalid", "Current password is incorrect");
}
@@ -36,17 +36,21 @@ public sealed class Chat : AggregateRoot<Guid>
public string? Avatar { get; private set; }
public DateTime CreatedAt { get; private set; }
public long LastMessageSequenceId { get; private set; }
public bool IsImporting { get; private set; }
public Guid? ImportJobId { get; private set; }
private readonly List<ChatMember> _members = new();
public IReadOnlyCollection<ChatMember> Members => _members.AsReadOnly();
private Chat(Guid id, ChatType type, string? name, string? avatar, string? description = null) : base(id)
private Chat(Guid id, ChatType type, string? name, string? avatar, string? description = null, bool isImporting = false, Guid? importJobId = null) : base(id)
{
Type = type;
Name = name;
Avatar = avatar;
Description = description;
CreatedAt = DateTime.UtcNow;
IsImporting = isImporting;
ImportJobId = importJobId;
}
public static Chat CreatePersonal()
@@ -63,13 +67,18 @@ public sealed class Chat : AggregateRoot<Guid>
return chat;
}
public static Chat Create(string? name, ChatType type, string? avatar = null, string? description = null)
public static Chat Create(string? name, ChatType type, string? avatar = null, string? description = null, bool isImporting = false, Guid? importJobId = null)
{
var chat = new Chat(Guid.NewGuid(), type, name, avatar, description);
var chat = new Chat(Guid.NewGuid(), type, name, avatar, description, isImporting, importJobId);
chat.RaiseDomainEvent(new ChatCreatedDomainEvent(chat));
return chat;
}
public void CompleteImport()
{
IsImporting = false;
}
public void AddMember(Guid userId, string role = "member")
{
if (_members.Any(m => m.UserId == userId))
@@ -0,0 +1,8 @@
namespace Knot.Contracts.Conversations.Domain;
public static class ChatConstants
{
public const int DefaultMessageQueryLimit = 50;
public const int MaxSharedMediaQueryLimit = 1000;
public const int MaxGroupNameLength = 100;
}
@@ -0,0 +1,19 @@
using Knot.Shared.Kernel;
namespace Knot.Contracts.Conversations.Domain;
public static class ChatErrors
{
public static readonly Error ChatNotFound = new Error("Chat.NotFound", "Чат не найден");
public static readonly Error OnlyOwnerCanUpdate = new Error("Chat.OnlyOwnerCanUpdate", "Только владелец может редактировать чат");
public static readonly Error Unauthorized = new Error("Chat.Unauthorized", "Нет доступа к этому чату");
public static readonly Error FoldersDisabled = new Error("Chat.FoldersDisabled", "Папки отключены");
public static readonly Error FileEmpty = new Error("Chat.FileEmpty", "Файл пуст");
public static Error FileTooLarge(long maxMb) => new Error("Chat.FileTooLarge", $"Файл слишком большой (максимум {maxMb} МБ)");
public static readonly Error ChatsNotFound = new Error("Chat.NotFound", "Чат не найден");
public static readonly Error ChatsForbidden = new Error("Chat.Forbidden", "Доступ запрещен");
public static readonly Error MediaDisabled = new Error("Chat.MediaDisabled", "Медиафайлы отключены");
public static readonly Error PollsDisabled = new Error("Chat.PollsDisabled", "Опросы отключены");
public static readonly Error NotFound = new Error("Chat.NotFound", "Не найдено");
public static readonly Error NotMember = new Error("Chat.NotMember", "Вы не являетесь участником чата");
}
@@ -3,4 +3,5 @@ namespace Knot.Contracts.Messaging.Application.Abstractions;
public interface IMessageNotifier
{
Task NotifyNewMessageAsync(Guid chatId, object messagePayload, CancellationToken cancellationToken);
Task NotifyMessageUpdateAsync(Guid chatId, string updateType, object updatePayload, CancellationToken cancellationToken);
}
@@ -8,9 +8,11 @@ public interface IMessageRepository
Task<Message?> GetByIdAsync(Guid id, CancellationToken cancellationToken);
Task<List<Message>> GetChatMessagesAsync(Guid chatId, int limit, int offset, CancellationToken cancellationToken);
Task<Message?> GetLatestChatMessageAsync(Guid chatId, CancellationToken cancellationToken);
Task<List<Message>> GetPinnedMessagesAsync(Guid chatId, CancellationToken cancellationToken);
Task<List<Message>> SearchMessagesAsync(string query, Guid? chatId, Guid requestingUserId, CancellationToken cancellationToken);
Task<List<Message>> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, int limit, CancellationToken cancellationToken);
Task<List<Message>> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, long? sequenceId, int limit, CancellationToken cancellationToken);
Task<List<Message>> GetChatMessagesAroundAsync(Guid chatId, long sequenceId, int limit, CancellationToken cancellationToken);
Task<Message?> GetLastStoryMessageAsync(Guid chatId, Guid storyId, CancellationToken cancellationToken);
Task UpdateAsync(Message message, CancellationToken cancellationToken);
@@ -21,7 +21,7 @@ public class MediaMessage : Message
public MediaMessage(Guid id, Guid chatId, Guid senderId, string mediaType, string? caption, Guid? replyToId, Guid? forwardedFromId, DateTime createdAt, bool isImported)
: this(id, chatId, senderId, Enum.TryParse<MediaType>(mediaType, true, out var mt) ? mt : MediaType.File, caption, replyToId, forwardedFromId, createdAt, isImported) { }
public void AddMedia(string type, string url, string? filename, long? size) => _media.Add(new Media { Type = type, Url = url, FileId = filename, Size = size });
public void AddMedia(string type, string url, string? filename, long? size, string? duration = null) => _media.Add(new Media { Type = type, Url = url, Filename = filename, FileId = filename, Size = size, Duration = duration });
public override void Edit(string newCaption) => base.Edit(newCaption);
@@ -7,26 +7,32 @@ public class PollMessage : Message
{
public override string Type => "poll";
public override string? Content { get; protected set; }
public List<PollOption> Options { get; } = new();
public List<PollVote> Votes { get; } = new();
public List<PollOption> Options { get; set; } = new();
public List<PollVote> Votes { get; set; } = new();
public bool IsMultipleChoice { get; set; }
public bool IsAnonymous { get; set; }
public DateTime? ExpiresAt { get; set; }
public bool IsClosed { get; set; }
public PollMessage() : base() { }
public PollMessage(Guid id, Guid chatId, Guid senderId, string? question, List<string>? options, bool isAnonymous, bool isMultiple, DateTime? expiresAt, Guid? replyToId, Guid? forwardedFromId, DateTime createdAt, bool isImported)
public PollMessage(Guid id, Guid chatId, Guid senderId, string? question, List<PollOption>? options, bool isAnonymous, bool isMultiple, DateTime? expiresAt, Guid? replyToId, Guid? forwardedFromId, DateTime createdAt, bool isImported)
: base(id, chatId, senderId, replyToId, forwardedFromId, createdAt, isImported)
{
Content = question ?? "Poll";
if (options != null)
{
foreach (var opt in options)
{
Options.Add(new PollOption { Text = opt });
}
}
Options = options ?? new List<PollOption>();
IsAnonymous = isAnonymous;
IsMultipleChoice = isMultiple;
ExpiresAt = expiresAt;
}
public static PollMessage Create(Guid id, Guid chatId, Guid senderId, string? question, List<string> options, bool isAnonymous, bool isMultiple, DateTime? expiresAt, Guid? replyToId, Guid? forwardedFromId)
{
var poll = new PollMessage(id, chatId, senderId, question, null, isAnonymous, isMultiple, expiresAt, replyToId, forwardedFromId, DateTime.UtcNow, false);
foreach (var opt in options)
{
poll.Options.Add(new PollOption { Text = opt });
}
return poll;
}
}
@@ -2,6 +2,7 @@ namespace Knot.Contracts.Messaging.Domain;
public class PollOption
{
public Guid Id { get; set; } = Guid.NewGuid();
public string Text { get; set; } = string.Empty;
public int VoteCount { get; set; }
}
@@ -4,7 +4,7 @@ namespace Knot.Contracts.Messaging.Domain;
public class PollVote
{
public Guid OptionIndex { get; set; }
public Guid OptionId { get; set; }
public Guid UserId { get; set; }
public DateTime VotedAt { get; set; }
}
@@ -9,4 +9,8 @@ public static class ProfilesErrors
public static Error AvatarNotFound => new("Profiles.AvatarNotFound", "Avatar not found");
public static Error InvalidAvatarFormat => new("Profiles.InvalidAvatarFormat", "Invalid avatar format");
public static Error AvatarUploadFailed => new("Profiles.AvatarUploadFailed", "Avatar upload failed");
public static Error BioTooLong => new("Profiles.BioTooLong", "Bio must be 200 characters or less");
public static Error StatusTextTooLong => new("Profiles.StatusTextTooLong", "Status text must be 50 characters or less");
public static Error StatusEmpty => new("Statuses.Empty", "Status emoji or text is required");
public static Error InvalidPreset => new("Statuses.InvalidPreset", "Unknown status preset");
}
@@ -0,0 +1,9 @@
namespace Knot.Contracts.Profiles.Application.DTOs;
public sealed class StatusPresetDto
{
public string Id { get; set; } = "";
public string Emoji { get; set; } = "";
public string TextRu { get; set; } = "";
public string TextEn { get; set; } = "";
}
@@ -12,5 +12,13 @@ public class UserProfileDto
public DateTime? LastSeen { get; set; }
public DateTime? Birthday { get; set; }
public bool IsPremium { get; set; }
public UserStatusDto? Status { get; set; }
public Guid? CurrentStatusId { get; set; }
public string? StatusText { get; set; }
public string? StatusEmoji { get; set; }
public DateTime? StatusExpiresAt { get; set; }
public bool IsInvisible { get; set; }
public DateTime CreatedAt { get; set; }
}
@@ -0,0 +1,11 @@
namespace Knot.Contracts.Profiles.Application.DTOs;
public sealed class UserStatusDto
{
public Guid Id { get; set; }
public string Type { get; set; } = "Custom";
public string Emoji { get; set; } = "";
public string Text { get; set; } = "";
public DateTime CreatedAt { get; set; }
public DateTime? ExpiresAt { get; set; }
}
@@ -0,0 +1,11 @@
using Knot.Contracts.Profiles.Application.DTOs;
using Knot.Shared.Kernel;
namespace Knot.Contracts.Profiles.Domain;
public interface IProfileStatusWriter
{
Task<Result<UserProfileDto>> SetCustomAsync(Guid userId, string emoji, string text, DateTime? expiresAt, string? presetKey, CancellationToken cancellationToken = default);
Task<Result<UserProfileDto>> ClearAsync(Guid userId, CancellationToken cancellationToken = default);
}
@@ -0,0 +1,14 @@
using Knot.Contracts.Profiles.Application.DTOs;
namespace Knot.Contracts.Profiles.Domain;
public interface IUserStatusRepository
{
Task<UserStatusDto?> GetByIdAsync(Guid id, CancellationToken cancellationToken = default);
Task<IReadOnlyDictionary<Guid, UserStatusDto>> GetByIdsAsync(IEnumerable<Guid> ids, CancellationToken cancellationToken = default);
Task InsertAsync(UserStatusDto dto, Guid userId, CancellationToken cancellationToken = default);
Task DeleteAsync(Guid id, CancellationToken cancellationToken = default);
}
+8 -2
View File
@@ -41,6 +41,7 @@ using MediatR;
JwtSecurityTokenHandler.DefaultInboundClaimTypeMap.Clear();
Encoding.RegisterProvider(CodePagesEncodingProvider.Instance);
var builder = WebApplication.CreateBuilder(args);
@@ -48,7 +49,7 @@ var builder = WebApplication.CreateBuilder(args);
// Маппинг стандартных переменных окружения в иерархию .NET
var envMappings = new Dictionary<string, string?>
{
["ConnectionStrings:DefaultConnection"] = builder.Configuration["DATABASE_URL"] ?? "Host=localhost;Database=knot;Username=postgres;Password=postgres",
["ConnectionStrings:DefaultConnection"] = builder.Configuration["DATABASE_URL"] ?? builder.Configuration.GetConnectionString("DefaultConnection") ?? "Host=localhost;Database=knot;Username=postgres;Password=postgres",
["ConnectionStrings:MongoConnection"] = builder.Configuration["MONGO_CONNECTION"],
["Jwt:Secret"] = builder.Configuration["JWT_SECRET"],
["Jwt:Issuer"] = builder.Configuration["JWT_ISSUER"],
@@ -81,6 +82,7 @@ builder.Services.AddStoriesModule(builder.Configuration);
builder.Services.AddKlipyModule();
builder.Services.AddAdminModule();
builder.Services.AddWebRtcModule();
builder.Services.AddTelegramImportModule();
builder.Services.AddSharedInfrastructure(builder.Configuration);
// CQRS / MediatR для команд в Host (например, AdminController)
@@ -93,7 +95,10 @@ builder.Services.AddMediatR(cfg => cfg.RegisterServicesFromAssemblies(
typeof(Knot.Modules.Stories.DependencyInjection).Assembly,
typeof(Knot.Modules.Klipy.DependencyInjection).Assembly,
typeof(Knot.Modules.Relations.DependencyInjection).Assembly,
typeof(Knot.Modules.WebRtc.DependencyInjection).Assembly
typeof(Knot.Modules.WebRtc.DependencyInjection).Assembly,
typeof(Knot.Modules.TelegramImport.DependencyInjection).Assembly,
typeof(Knot.Modules.Auth.Infrastructure.Persistence.AuthDbContext).Assembly,
typeof(Knot.Modules.Profiles.DependencyInjection).Assembly
));
// Настройка CORS
@@ -245,6 +250,7 @@ app.MapStoriesEndpoints();
app.MapContactsEndpoints();
app.MapSettingsEndpoints();
app.MapProfilesEndpoints();
app.MapStatusesEndpoints();
app.MapFederationEndpoints();
app.MapKlipyEndpoints();
app.MapChatsEndpoints();
@@ -8,7 +8,7 @@ using Knot.Contracts.Auth.Infrastructure.Persistence;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Infrastructure.Persistence;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Storage.Abstractions;
using Knot.Shared.Kernel.Storage;
using Knot.Contracts.Stories.Infrastructure.Persistence;
using Knot.Modules.Admin.Application.Admin.DTOs;
using MediatR;
@@ -51,7 +51,7 @@ internal sealed class CleanRunCommandHandler : ICommandHandler<CleanRunCommand,
var orphanMessages = await _messageQueryService.GetOrphanedMessagesAsync(activeChatIds, cancellationToken);
var keptMessages = allMessages
.Where(m => !orphanMessages.Any(om => om.Id == m.Id))
.Where(m => !orphanMessages.Any(om => om.Id == m.Id) && !m.IsDeleted)
.ToList();
var allMinioFiles = (await _fileStorage.ListFilesAsync()).ToList();
@@ -1,4 +1,4 @@
using System;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
@@ -8,6 +8,7 @@ using Knot.Contracts.Auth.Infrastructure.Persistence;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Infrastructure.Persistence;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Shared.Kernel.Storage;
using Knot.Contracts.Stories.Infrastructure.Persistence;
using Knot.Modules.Admin.Application.Admin.DTOs;
using MediatR;
@@ -23,17 +24,20 @@ internal sealed class CleanDryRunQueryHandler : IQueryHandler<CleanDryRunQuery,
private readonly IAuthDbContext _authDbContext;
private readonly IChatsDbContext _chatsDbContext;
private readonly IStoryCollection _storyCollection;
private readonly IFileStorageService _fileStorage;
public CleanDryRunQueryHandler(
Knot.Contracts.Messaging.Application.Abstractions.IMessageQueryService messageService,
IAuthDbContext authDbContext,
IChatsDbContext chatsDbContext,
IStoryCollection storyCollection)
IStoryCollection storyCollection,
IFileStorageService fileStorage)
{
_messageService = messageService;
_authDbContext = authDbContext;
_chatsDbContext = chatsDbContext;
_storyCollection = storyCollection;
_fileStorage = fileStorage;
}
public async Task<Result<CleanDryRunResult>> Handle(CleanDryRunQuery request, CancellationToken ct)
@@ -46,25 +50,42 @@ internal sealed class CleanDryRunQueryHandler : IQueryHandler<CleanDryRunQuery,
var orphanedMediaCount = orphanedMessages.Count(m => m.MediaUrl != null);
var orphanedMessageCount = orphanedMessages.Count;
var validIds = new HashSet<string>();
foreach (var msg in orphanedMessages.Where(m => m.MediaUrl != null))
{
var parts = msg.MediaUrl.Split('/');
var fileId = parts.LastOrDefault();
if (!string.IsNullOrEmpty(fileId))
{
validIds.Add(fileId);
}
}
var allMessages = await _messageService.GetAllMessagesAsync(ct);
var keptMessages = allMessages
.Where(m => !orphanedMessages.Any(om => om.Id == m.Id) && !m.IsDeleted)
.ToList();
var allUsers = await _authDbContext.Users.ToListAsync(ct);
var stories = await _storyCollection.GetAllAsync(ct);
var validUrls = new HashSet<string>();
var activeMessageUrls = keptMessages.Where(m => m.Media != null).SelectMany(m => m.Media!).Select(x => x.Url).Where(u => !string.IsNullOrEmpty(u));
var activeChatUrls = _chatsDbContext.Chats.Select(c => c.Avatar).Where(u => !string.IsNullOrEmpty(u));
var activeUserUrls = allUsers.Select(u => u.Avatar).Where(u => !string.IsNullOrEmpty(u));
var activeStoryUrls = stories.Select(s => s.MediaUrl).Where(u => !string.IsNullOrEmpty(u));
foreach (var u in activeMessageUrls) validUrls.Add(u!);
foreach (var u in activeChatUrls) validUrls.Add(u!);
foreach (var u in activeUserUrls) validUrls.Add(u!);
foreach (var u in activeStoryUrls) validUrls.Add(u!);
var validFileIds = validUrls
.Where(u => u.Contains("/api/files/"))
.Select(u => u.Split('/').Last())
.ToHashSet();
var allMinioFiles = await _fileStorage.ListFilesAsync();
long orphanedFileSize = allMinioFiles
.Where(f => !validFileIds.Contains(f.FileId))
.Sum(f => f.Size);
var expiredStoriesCount = stories.Count(s => s.ExpiresAt.HasValue && s.ExpiresAt.Value < DateTime.UtcNow);
var expiredStoriesSize = stories.Where(s => s.ExpiresAt.HasValue && s.ExpiresAt.Value < DateTime.UtcNow).Sum(s => s.MediaUrl?.Length ?? 0);
return Result.Success(new CleanDryRunResult(
orphanedMessageCount,
orphanedMediaCount,
0,
orphanedFileSize,
expiredStoriesCount,
expiredStoriesSize
));
@@ -131,6 +131,17 @@ public static class AdminEndpoints
return Results.Ok(result.Value);
});
group.MapPost("clean/run", async (ISender sender, CancellationToken ct) =>
{
var result = await sender.Send(new CleanRunCommand(), ct);
if (!result.IsSuccess)
{
Console.WriteLine($"[Admin] Cleanup Run Error: {result.Error.Description}");
return Results.BadRequest(new { error = result.Error.Description });
}
return Results.Ok(result.Value);
});
group.MapGet("timezones", () =>
{
// Получаем все системные часовые пояса и формируем удобный для фронтенда формат
@@ -0,0 +1,48 @@
using BCrypt.Net;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Contracts.Auth.Application.Auth.DTOs;
using Knot.Modules.Auth.Domain;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
namespace Knot.Modules.Auth.Application.Users.ChangePassword;
public sealed record ChangePasswordCommand(
Guid UserId,
string OldPassword,
string NewPassword,
string ConfirmPassword) : ICommand;
internal sealed class ChangePasswordCommandHandler : ICommandHandler<ChangePasswordCommand>
{
private readonly IAuthDbContext _dbContext;
public ChangePasswordCommandHandler(IAuthDbContext dbContext)
{
_dbContext = dbContext;
}
public async Task<Result> Handle(ChangePasswordCommand request, CancellationToken cancellationToken)
{
if (request.NewPassword != request.ConfirmPassword)
return Result.Failure(AuthErrors.PasswordConfirmationMismatch);
if (request.NewPassword.Length < 8)
return Result.Failure(AuthErrors.PasswordTooShort);
var user = await _dbContext.Set<User>()
.FirstOrDefaultAsync(u => u.Id == request.UserId, cancellationToken);
if (user is null)
return Result.Failure(AuthErrors.UserNotFound);
if (!BCrypt.Net.BCrypt.Verify(request.OldPassword, user.PasswordHash))
return Result.Failure(AuthErrors.OldPasswordInvalid);
user.ChangePassword(BCrypt.Net.BCrypt.HashPassword(request.NewPassword));
user.SetRefreshToken(null);
await _dbContext.SaveChangesAsync(cancellationToken);
return Result.Success();
}
}
@@ -0,0 +1,106 @@
// <auto-generated />
using System;
using Knot.Modules.Auth.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
#nullable disable
namespace Knot.Modules.Auth.Migrations
{
[DbContext(typeof(AuthDbContext))]
[Migration("20260407181656_AddUserInfoFields")]
partial class AddUserInfoFields
{
/// <inheritdoc />
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasDefaultSchema("identity")
.HasAnnotation("ProductVersion", "10.0.4")
.HasAnnotation("Relational:MaxIdentifierLength", 63);
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Modules.Auth.Domain.User", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Avatar")
.HasColumnType("text");
b.Property<DateTime?>("BannedUntil")
.HasColumnType("timestamp with time zone");
b.Property<string>("Bio")
.HasColumnType("text");
b.Property<DateTime?>("Birthday")
.HasColumnType("timestamp with time zone");
b.Property<DateTime>("CreatedAt")
.HasColumnType("timestamp with time zone");
b.Property<string>("DisplayName")
.IsRequired()
.HasColumnType("text");
b.Property<string>("Domain")
.HasColumnType("text");
b.Property<string>("Email")
.HasColumnType("text");
b.Property<bool>("HideStatus")
.HasColumnType("boolean");
b.Property<bool>("HideStoryViews")
.HasColumnType("boolean");
b.Property<bool>("IsBanned")
.HasColumnType("boolean");
b.Property<bool>("IsExternal")
.HasColumnType("boolean");
b.Property<bool>("IsOnline")
.HasColumnType("boolean");
b.Property<DateTime?>("LastSeen")
.HasColumnType("timestamp with time zone");
b.Property<string>("PasswordHash")
.IsRequired()
.HasColumnType("text");
b.Property<string>("PhoneNumber")
.HasColumnType("text");
b.Property<string>("RefreshToken")
.HasColumnType("text");
b.Property<string>("UserDomain")
.HasColumnType("text");
b.Property<string>("Username")
.IsRequired()
.HasMaxLength(50)
.HasColumnType("character varying(50)");
b.HasKey("Id");
b.HasIndex("Username")
.IsUnique();
b.ToTable("Users", "identity");
});
#pragma warning restore 612, 618
}
}
}
@@ -0,0 +1,37 @@
using System;
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace Knot.Modules.Auth.Migrations
{
/// <inheritdoc />
public partial class AddUserInfoFields : Migration
{
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.Sql(@"
DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema='identity' AND table_name='Users' AND column_name='BannedUntil') THEN
ALTER TABLE identity.""Users"" ADD ""BannedUntil"" timestamp with time zone;
END IF;
IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema='identity' AND table_name='Users' AND column_name='PhoneNumber') THEN
ALTER TABLE identity.""Users"" ADD ""PhoneNumber"" text;
END IF;
IF NOT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema='identity' AND table_name='Users' AND column_name='RefreshToken') THEN
ALTER TABLE identity.""Users"" ADD ""RefreshToken"" text;
END IF;
END $$;");
}
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropColumn(name: "BannedUntil", schema: "identity", table: "Users");
migrationBuilder.DropColumn(name: "PhoneNumber", schema: "identity", table: "Users");
migrationBuilder.DropColumn(name: "RefreshToken", schema: "identity", table: "Users");
}
}
}
@@ -1,4 +1,4 @@
// <auto-generated />
// <auto-generated />
using System;
using Knot.Modules.Auth.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
@@ -23,7 +23,7 @@ namespace Knot.Modules.Auth.Migrations
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Contracts.Auth.Domain.User", b =>
modelBuilder.Entity("Knot.Modules.Auth.Domain.User", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
@@ -32,6 +32,9 @@ namespace Knot.Modules.Auth.Migrations
b.Property<string>("Avatar")
.HasColumnType("text");
b.Property<DateTime?>("BannedUntil")
.HasColumnType("timestamp with time zone");
b.Property<string>("Bio")
.HasColumnType("text");
@@ -57,10 +60,10 @@ namespace Knot.Modules.Auth.Migrations
b.Property<bool>("HideStoryViews")
.HasColumnType("boolean");
b.Property<bool>("IsExternal")
b.Property<bool>("IsBanned")
.HasColumnType("boolean");
b.Property<bool>("IsBanned")
b.Property<bool>("IsExternal")
.HasColumnType("boolean");
b.Property<bool>("IsOnline")
@@ -73,6 +76,15 @@ namespace Knot.Modules.Auth.Migrations
.IsRequired()
.HasColumnType("text");
b.Property<string>("PhoneNumber")
.HasColumnType("text");
b.Property<string>("RefreshToken")
.HasColumnType("text");
b.Property<string>("UserDomain")
.HasColumnType("text");
b.Property<string>("Username")
.IsRequired()
.HasMaxLength(50)
@@ -2,6 +2,7 @@ using Knot.Shared.Kernel;
using Knot.Modules.Auth.Application.Users.Login;
using Knot.Modules.Auth.Application.Users.Register;
using Knot.Modules.Auth.Application.Users.GetMe;
using Knot.Modules.Auth.Application.Users.ChangePassword;
using MediatR;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Http;
@@ -12,6 +13,8 @@ namespace Knot.Modules.Auth.Presentation.Endpoints;
public static class AuthEndpoints
{
public sealed record ChangePasswordRequest(string OldPassword, string NewPassword, string ConfirmPassword);
public static void MapAuthEndpoints(this WebApplication app)
{
var group = app.MapGroup("api/auth");
@@ -33,5 +36,19 @@ public static class AuthEndpoints
var result = await sender.Send(new GetMeQuery(userContext.UserId), ct);
return result.IsSuccess ? Results.Ok(result.Value) : Results.NotFound();
}).RequireAuthorization();
group.MapPost("change-password", async ([FromBody] ChangePasswordRequest request, ISender sender, IUserContext userContext, CancellationToken ct) =>
{
var result = await sender.Send(
new ChangePasswordCommand(userContext.UserId, request.OldPassword, request.NewPassword, request.ConfirmPassword),
ct);
if (result.IsSuccess) return Results.Ok();
if (result.Error.Code == "Auth.OldPasswordInvalid")
return Results.StatusCode(StatusCodes.Status403Forbidden);
return Results.BadRequest(new { error = result.Error.Code ?? result.Error.Description });
}).RequireAuthorization();
}
}
@@ -1,11 +0,0 @@
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Application.Abstractions;
/// <summary>
/// Unit of Work специфичный для модуля Chats.
/// </summary>
public interface IChatsUnitOfWork : IUnitOfWork
{
}
@@ -1,11 +0,0 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Application.Abstractions;
public interface IUserDeleterService
{
Task<Result> DeleteUserAsync(Guid userId, CancellationToken cancellationToken);
}
@@ -1,8 +0,0 @@
using System;
namespace Knot.Modules.Conversations.Application.Abstractions;
public interface IUserStatusService
{
bool IsUserOnline(string userId);
}
@@ -4,8 +4,8 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Storage;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
using SixLabors.ImageSharp;
@@ -2,8 +2,8 @@ using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
@@ -1,7 +1,7 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
namespace Knot.Modules.Conversations.Application.Chats.Create;
@@ -5,9 +5,9 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
@@ -63,6 +63,9 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
}
}
var pinnedMessages = await _messageRepository.GetPinnedMessagesAsync(chat.Id, cancellationToken);
foreach (var pm in pinnedMessages) userIdsToFetch.Add(pm.SenderId);
var usersInfo = await _userProvider.GetUsersInfoAsync(userIdsToFetch, cancellationToken);
var members = new List<ChatMemberDto>();
@@ -88,55 +91,19 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
var messagesList = new List<ChatMessageDto>();
if (latestMessage != null)
{
usersInfo.TryGetValue(latestMessage.SenderId, out var senderObj);
var reactionsWithUser = new List<ReactionDto>();
foreach (var reaction in latestReactions)
{
usersInfo.TryGetValue(reaction.UserId, out var reactionUser);
reactionsWithUser.Add(new ReactionDto(
reaction.Id,
reaction.Emoji,
reaction.UserId,
reactionUser != null
? new MessageSenderDto(reactionUser.Id, reactionUser.Username, reactionUser.DisplayName, reactionUser.Avatar)
: new MessageSenderDto(reaction.UserId, "unknown", "Unknown", null)
));
messagesList.Add(MessageMapper.MapToDto(
latestMessage,
usersInfo,
latestReactions,
chat.Members.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId).Select(m => m.UserId)));
}
var readByList = chat.Members
.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId)
.Select(m => new ReadByDto(m.UserId))
.ToList();
var textMessage = latestMessage as TextMessage;
var mediaMessage = latestMessage as MediaMessage;
var storyMessage = latestMessage as StoryMessage;
messagesList.Add(new ChatMessageDto(
latestMessage.Id,
latestMessage.ChatId,
latestMessage.SenderId,
latestMessage.Content,
latestMessage.Type,
latestMessage.ReplyToId,
textMessage?.Quote,
storyMessage?.StoryId,
storyMessage?.StoryMediaUrl,
storyMessage?.StoryMediaType,
latestMessage.IsEdited,
latestMessage.IsDeleted,
latestMessage.CreatedAt,
latestMessage.SequenceId,
mediaMessage?.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList() ?? new List<MediaDto>(),
senderObj != null ? new MessageSenderDto(
senderObj.Id,
senderObj.Username,
senderObj.DisplayName,
senderObj.Avatar
) : new MessageSenderDto(latestMessage.SenderId, "unknown", "Unknown", null),
reactionsWithUser,
readByList
var pinnedDtoList = new List<PinnedMessageDto>();
foreach (var pm in pinnedMessages)
{
pinnedDtoList.Add(new PinnedMessageDto(
pm.Id,
MessageMapper.MapToDto(pm, usersInfo, new List<MessageReaction>(), new List<Guid>())
));
}
@@ -152,6 +119,7 @@ internal sealed class GetChatByIdQueryHandler : IQueryHandler<GetChatByIdQuery,
chat.CreatedAt,
members,
messagesList,
pinnedDtoList,
unreadCount
);
@@ -5,9 +5,9 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
@@ -59,6 +59,9 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
}
}
var pinnedMessages = await _messageRepository.GetPinnedMessagesAsync(chat.Id, cancellationToken);
foreach (var pm in pinnedMessages) userIdsToFetch.Add(pm.SenderId);
var usersInfo = await _userProvider.GetUsersInfoAsync(userIdsToFetch, cancellationToken);
var members = new List<ChatMemberDto>();
@@ -85,53 +88,19 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
if (latestMessage != null)
{
usersInfo.TryGetValue(latestMessage.SenderId, out var senderObj);
var reactionsWithUser = new List<ReactionDto>();
foreach (var reaction in latestReactions)
{
usersInfo.TryGetValue(reaction.UserId, out var reactionUser);
reactionsWithUser.Add(new ReactionDto(
reaction.Id,
reaction.Emoji,
reaction.UserId,
reactionUser != null
? new MessageSenderDto(reactionUser.Id, reactionUser.Username, reactionUser.DisplayName, reactionUser.Avatar)
: new MessageSenderDto(reaction.UserId, "unknown", "Unknown", null)
));
messagesList.Add(MessageMapper.MapToDto(
latestMessage,
usersInfo,
latestReactions,
chat.Members.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId).Select(m => m.UserId)));
}
var textMessage = latestMessage as TextMessage;
var mediaMessage = latestMessage as MediaMessage;
var storyMessage = latestMessage as StoryMessage;
messagesList.Add(new ChatMessageDto(
latestMessage.Id,
latestMessage.ChatId,
latestMessage.SenderId,
latestMessage.Content,
latestMessage.Type,
latestMessage.ReplyToId,
textMessage?.Quote,
storyMessage?.StoryId,
storyMessage?.StoryMediaUrl,
storyMessage?.StoryMediaType,
latestMessage.IsEdited,
latestMessage.IsDeleted,
latestMessage.CreatedAt,
latestMessage.SequenceId,
mediaMessage?.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList() ?? new List<MediaDto>(),
senderObj != null ? new MessageSenderDto(
senderObj.Id,
senderObj.Username,
senderObj.DisplayName,
senderObj.Avatar
) : new MessageSenderDto(latestMessage.SenderId, "unknown", "Unknown", null),
reactionsWithUser,
chat.Members.Where(m => m.LastReadSequenceId >= latestMessage.SequenceId && m.UserId != latestMessage.SenderId).Select(m => new ReadByDto(m.UserId)).ToList(),
(latestMessage as CallMessage)?.CallType,
(latestMessage as CallMessage)?.CallStatus,
(latestMessage as CallMessage)?.Duration
var pinnedDtoList = new List<PinnedMessageDto>();
foreach (var pm in pinnedMessages)
{
pinnedDtoList.Add(new PinnedMessageDto(
pm.Id,
MessageMapper.MapToDto(pm, usersInfo, new List<MessageReaction>(), new List<Guid>())
));
}
@@ -147,7 +116,10 @@ internal sealed class GetChatsQueryHandler : IQueryHandler<GetChatsQuery, List<C
chat.CreatedAt,
members,
messagesList,
unreadCount
pinnedDtoList,
unreadCount,
chat.IsImporting,
chat.ImportJobId
));
}
@@ -1,6 +1,6 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
namespace Knot.Modules.Conversations.Application.Chats.GetOrCreateFavorites;
@@ -1,12 +1,13 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Shared.Kernel.Storage;
namespace Knot.Modules.Conversations.Application.Chats.LeaveOrDelete;
@@ -15,11 +16,19 @@ public record LeaveOrDeleteChatCommand(Guid ChatId, Guid UserId) : ICommand<Succ
internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrDeleteChatCommand, SuccessResponse>
{
private readonly IChatRepository _chatRepository;
private readonly IMessageRepository _messageRepository;
private readonly IFileStorageService _fileStorage;
private readonly IChatsUnitOfWork _uow;
public LeaveOrDeleteChatCommandHandler(IChatRepository chatRepository, IChatsUnitOfWork uow)
public LeaveOrDeleteChatCommandHandler(
IChatRepository chatRepository,
IMessageRepository messageRepository,
IFileStorageService fileStorage,
IChatsUnitOfWork uow)
{
_chatRepository = chatRepository;
_messageRepository = messageRepository;
_fileStorage = fileStorage;
_uow = uow;
}
@@ -36,13 +45,18 @@ internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrD
return Result.Failure<SuccessResponse>(ChatErrors.Unauthorized);
}
if (chat.Type == ChatType.Group)
// If it's a private chat or the last member leaving a group, delete everything
bool shouldDeleteEverything = chat.Type != ChatType.Group || chat.Members.Count <= 1;
if (chat.Type == ChatType.Group && !shouldDeleteEverything)
{
chat.RemoveMember(request.UserId);
_chatRepository.Update(chat);
}
else
{
// DELETE ALL MESSAGES AND FILES FIRST
await DeleteChatMediaAndMessagesAsync(chat.Id, cancellationToken);
_chatRepository.Remove(chat);
}
@@ -50,5 +64,45 @@ internal sealed class LeaveOrDeleteChatCommandHandler : ICommandHandler<LeaveOrD
return Result.Success(new SuccessResponse(true));
}
}
private async Task DeleteChatMediaAndMessagesAsync(Guid chatId, CancellationToken ct)
{
try
{
// Get all messages directly from Mongo (not paged)
var messages = await _messageRepository.GetChatMessagesAsync(chatId, int.MaxValue, 0, ct);
foreach (var msg in messages)
{
if (msg is Knot.Contracts.Messaging.Domain.MediaMessage mediaMsg)
{
foreach (var media in mediaMsg.Media)
{
if (!string.IsNullOrEmpty(media.Url))
{
var fileId = ExtractFileId(media.Url);
if (!string.IsNullOrEmpty(fileId))
{
await _fileStorage.DeleteFileAsync(fileId);
}
}
}
}
}
await _messageRepository.DeleteChatMessagesAsync(chatId, ct);
}
catch (Exception ex)
{
// Log if possible, but don't fail chat deletion
Console.WriteLine($"[Cleanup] Error deleting chat media: {ex.Message}");
}
}
private string? ExtractFileId(string url)
{
var lastSlash = url.LastIndexOf('/');
if (lastSlash == -1) return null;
var id = url[(lastSlash + 1)..];
if (id.Contains('?')) id = id[..id.IndexOf('?')];
return id;
}
}
@@ -3,8 +3,8 @@ using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
@@ -2,8 +2,8 @@ using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Application.DTOs;
using MediatR;
using System.Linq;
@@ -2,8 +2,8 @@ using System;
using System.Threading;
using System.Threading.Tasks;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using MediatR;
using System.Linq;
@@ -12,6 +12,9 @@ public record ChatDto(
DateTime CreatedAt,
List<ChatMemberDto> Members,
List<ChatMessageDto> Messages,
int UnreadCount
List<PinnedMessageDto> PinnedMessages,
int UnreadCount,
bool IsImporting = false,
Guid? ImportJobId = null
);
@@ -1,4 +1,4 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using System;
namespace Knot.Modules.Conversations.Application.DTOs;
@@ -24,6 +24,14 @@ public record ChatMessageDto(
List<ReadByDto> ReadBy,
string? CallType = null,
string? CallStatus = null,
int? Duration = null
int? Duration = null,
List<PollOptionDto>? PollOptions = null,
bool? PollIsMultipleChoice = null,
bool? PollIsAnonymous = null,
bool? PollIsClosed = null,
List<Guid>? UserVotedOptionIds = null
);
public record PollOptionDto(Guid Id, string Text, int VoteCount, List<MessageSenderDto>? Voters = null, List<Guid>? VoterIds = null);
@@ -1,6 +1,6 @@
using System;
using System.Collections.Generic;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
namespace Knot.Modules.Conversations.Application.DTOs;
@@ -7,6 +7,6 @@ public record MediaDto(
string Type,
string? Url,
string? Filename,
long? Size
long? Size,
string? Duration = null
);
@@ -27,7 +27,12 @@ public record MessageDetailDto(
List<MessageReactionDto> Reactions,
string? CallType = null,
string? CallStatus = null,
int? Duration = null
int? Duration = null,
List<PollOptionDto>? PollOptions = null,
bool? PollIsMultipleChoice = null,
bool? PollIsAnonymous = null,
bool? PollIsClosed = null,
List<Guid>? UserVotedOptionIds = null
);
public record ReplyToMessageDto(
@@ -0,0 +1,84 @@
using Knot.Shared.Kernel;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.DTOs;
namespace Knot.Modules.Conversations.Application.DTOs;
public static class MessageMapper
{
public static ChatMessageDto MapToDto(
Message message,
IReadOnlyDictionary<Guid, UserInfo> usersInfo,
IEnumerable<MessageReaction> reactions,
IEnumerable<Guid> readByUsers,
Guid? currentUserId = null)
{
usersInfo.TryGetValue(message.SenderId, out var senderObj);
var reactionsWithUser = new List<ReactionDto>();
foreach (var reaction in reactions)
{
usersInfo.TryGetValue(reaction.UserId, out var reactionUser);
reactionsWithUser.Add(new ReactionDto(
reaction.Id,
reaction.Emoji,
reaction.UserId,
reactionUser != null
? new MessageSenderDto(reactionUser.Id, reactionUser.Username, reactionUser.DisplayName, reactionUser.Avatar)
: new MessageSenderDto(reaction.UserId, "unknown", "Unknown", null)
));
}
var textMessage = message as TextMessage;
var mediaMessage = message as MediaMessage;
var storyMessage = message as StoryMessage;
var callMessage = message as CallMessage;
return new ChatMessageDto(
message.Id,
message.ChatId,
message.SenderId,
message.Content,
message.Type,
message.ReplyToId,
textMessage?.Quote,
storyMessage?.StoryId,
storyMessage?.StoryMediaUrl,
storyMessage?.StoryMediaType,
message.IsEdited,
message.IsDeleted,
message.CreatedAt,
message.SequenceId,
mediaMessage?.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size, media.Duration)).ToList() ?? new List<MediaDto>(),
senderObj != null ? new MessageSenderDto(
senderObj.Id,
senderObj.Username,
senderObj.DisplayName,
senderObj.Avatar
) : new MessageSenderDto(message.SenderId, "unknown", "Unknown", null),
reactionsWithUser,
readByUsers.Select(id => new ReadByDto(id)).ToList(),
callMessage?.CallType,
callMessage?.CallStatus,
callMessage?.Duration,
message is PollMessage pm ? pm.Options.Select(o => {
var voters = pm.IsAnonymous == false
? pm.Votes
.Where(v => v.OptionId == o.Id)
.Select(v => {
usersInfo.TryGetValue(v.UserId, out var vu);
return vu != null
? new MessageSenderDto(vu.Id, vu.Username, vu.DisplayName, vu.Avatar)
: new MessageSenderDto(v.UserId, "unknown", "Unknown", null);
})
.ToList()
: null;
return new PollOptionDto(o.Id, o.Text, o.VoteCount, voters, pm.IsAnonymous == false ? pm.Votes.Where(v => v.OptionId == o.Id).Select(v => v.UserId).ToList() : null);
}).ToList() : null,
(message as PollMessage)?.IsMultipleChoice,
(message as PollMessage)?.IsAnonymous,
(message as PollMessage)?.IsClosed,
(message is PollMessage poll && currentUserId.HasValue) ? poll.Votes.Where(v => v.UserId == currentUserId.Value).Select(v => v.OptionId).ToList() : null
);
}
}
@@ -0,0 +1,6 @@
namespace Knot.Modules.Conversations.Application.DTOs;
public record PinnedMessageDto(
Guid Id,
ChatMessageDto Message
);
@@ -1,5 +1,5 @@
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Settings.Application.Abstractions;
using Knot.Shared.Kernel;
using MediatR;
@@ -1,5 +1,5 @@
using Knot.Modules.Conversations.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Shared.Kernel;
using MediatR;
@@ -1,5 +1,5 @@
using global::Knot.Modules.Conversations.Application.Abstractions;
using global::Knot.Modules.Conversations.Domain;
using global::Knot.Contracts.Conversations.Application.Abstractions;
using global::Knot.Contracts.Conversations.Domain;
using global::Knot.Modules.Conversations.Infrastructure.SignalR;
using global::Knot.Shared.Kernel;
using MediatR;
@@ -0,0 +1,69 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.SignalR;
using Knot.Shared.Kernel;
using MediatR;
using Microsoft.AspNetCore.SignalR;
namespace Knot.Modules.Conversations.Application.Messages.Edit;
public sealed record EditMessageCommand(
Guid MessageId,
Guid ChatId,
Guid UserId,
string Content) : ICommand;
public sealed class EditMessageCommandHandler : ICommandHandler<EditMessageCommand>
{
private readonly IMessageRepository _messageRepository;
private readonly IChatsUnitOfWork _unitOfWork;
private readonly IHubContext<ChatHub> _hubContext;
public EditMessageCommandHandler(
IMessageRepository messageRepository,
IChatsUnitOfWork unitOfWork,
IHubContext<ChatHub> hubContext)
{
_messageRepository = messageRepository;
_unitOfWork = unitOfWork;
_hubContext = hubContext;
}
public async Task<Result> Handle(EditMessageCommand request, CancellationToken cancellationToken)
{
var message = await _messageRepository.GetByIdAsync(request.MessageId, cancellationToken);
if (message is null)
{
return Result.Failure(new Error("Message.NotFound", "Message not found."));
}
if (message.SenderId != request.UserId)
{
return Result.Failure(new Error("Message.Forbidden", "You can only edit your own messages."));
}
if (message.ChatId != request.ChatId)
{
return Result.Failure(new Error("Message.InvalidChat", "Message does not belong to this chat."));
}
message.Edit(request.Content);
await _messageRepository.UpdateAsync(message, cancellationToken);
await _unitOfWork.SaveChangesAsync(cancellationToken);
// Notify clients
await _hubContext.Clients.Group(request.ChatId.ToString()).SendAsync("message_edited", new
{
messageId = message.Id,
chatId = message.ChatId,
content = message.Content,
isEdited = true
});
return Result.Success();
}
}
@@ -5,15 +5,15 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Conversations.Application.Messages.GetMessages;
public record GetMessagesQuery(Guid UserId, Guid ChatId, string? Cursor) : IQuery<List<MessageDetailDto>>;
public record GetMessagesQuery(Guid UserId, Guid ChatId, string? Cursor, long? Pivot = null, int? Limit = null) : IQuery<List<MessageDetailDto>>;
internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery, List<MessageDetailDto>>
{
@@ -38,15 +38,34 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
return Result.Failure<List<MessageDetailDto>>(ChatErrors.ChatsForbidden);
}
List<Message> messages;
int queryLimit = request.Limit ?? ChatConstants.DefaultMessageQueryLimit;
if (request.Pivot.HasValue)
{
messages = await _messageRepository.GetChatMessagesAroundAsync(request.ChatId, request.Pivot.Value, queryLimit, cancellationToken);
}
else
{
DateTime? cursorDate = null;
if (!string.IsNullOrEmpty(request.Cursor) && DateTime.TryParse(request.Cursor, null, System.Globalization.DateTimeStyles.RoundtripKind, out var parsed))
long? cursorSequenceId = null;
if (!string.IsNullOrEmpty(request.Cursor))
{
if (long.TryParse(request.Cursor, out var seqId))
{
cursorSequenceId = seqId;
}
else if (DateTime.TryParse(request.Cursor, null, System.Globalization.DateTimeStyles.RoundtripKind, out var parsed))
{
cursorDate = parsed.ToUniversalTime();
}
}
messages = await _messageRepository.GetChatMessagesCursorAsync(request.ChatId, cursorDate, cursorSequenceId, queryLimit, cancellationToken);
}
var messages = await _messageRepository.GetChatMessagesCursorAsync(request.ChatId, cursorDate, ChatConstants.DefaultMessageQueryLimit, cancellationToken);
var result = new List<MessageDetailDto>();
var userIdsToFetch = new HashSet<Guid>();
var replyMessages = new Dictionary<Guid, Message>();
@@ -57,6 +76,14 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
{
userIdsToFetch.Add(m.SenderId);
if (m is PollMessage poll && !poll.IsAnonymous)
{
foreach (var vote in poll.Votes)
{
userIdsToFetch.Add(vote.UserId);
}
}
if (!m.ReplyToId.HasValue)
{
continue;
@@ -85,73 +112,77 @@ internal sealed class GetMessagesQueryHandler : IQueryHandler<GetMessagesQuery,
continue;
}
ReplyToMessageDto? replyToObj = null;
if (message.ReplyToId.HasValue && replyMessages.TryGetValue(message.ReplyToId.Value, out var replyMsg))
{
var senderObj = senders.TryGetValue(replyMsg.SenderId, out var rs)
? new MessageSenderDto(rs.Id, rs.Username, rs.DisplayName, rs.Avatar)
: null;
senders.TryGetValue(message.SenderId, out var sender);
reactionsByMessage.TryGetValue(message.Id, out var reactions);
replyToObj = new ReplyToMessageDto(
replyMsg.Id,
replyMsg.Content,
replyMsg.IsDeleted,
(replyMsg as MediaMessage)?.Media.Select(rm => new MediaDto(rm.Id, rm.Type, rm.Url, rm.Filename, rm.Size)).ToList() ?? new List<MediaDto>(),
senderObj
);
Message? replyMsg = null;
if (message.ReplyToId.HasValue)
{
replyMessages.TryGetValue(message.ReplyToId.Value, out replyMsg);
}
var reactionsWithUser = new List<MessageReactionDto>();
var messageReactions = reactionsByMessage.TryGetValue(message.Id, out var mr) ? mr : new List<MessageReaction>();
foreach (var reaction in messageReactions)
UserInfo? replySender = null;
if (replyMsg != null)
{
var userObj = senders.TryGetValue(reaction.UserId, out var reactionUser)
? new MessageSenderDto(reactionUser.Id, reactionUser.Username, reactionUser.DisplayName, reactionUser.Avatar)
: new MessageSenderDto(reaction.UserId, "unknown", "Unknown", null);
reactionsWithUser.Add(new MessageReactionDto(
reaction.Id,
reaction.Emoji,
reaction.UserId,
userObj
));
senders.TryGetValue(replyMsg.SenderId, out replySender);
}
var textMessage = message as TextMessage;
var mediaMessage = message as MediaMessage;
var storyMessage = message as StoryMessage;
result.Add(new MessageDetailDto(
message.Id,
message.ChatId,
message.SenderId,
message.Content,
message.Type,
message.Type.ToLower(),
message.ReplyToId,
replyToObj,
textMessage?.Quote,
replyMsg != null ? new ReplyToMessageDto(
replyMsg.Id,
replyMsg.Content,
replyMsg.IsDeleted,
replyMsg is MediaMessage mm ? mm.Media.Select(m => new MediaDto(m.Id, m.Type, m.Url, m.Filename, m.Size, m.Duration)).ToList() : new List<MediaDto>(),
replySender != null ? new MessageSenderDto(replySender.Id, replySender.Username, replySender.DisplayName, replySender.Avatar) : null
) : null,
message is TextMessage tm ? tm.Quote : null,
message.IsEdited,
message.IsDeleted,
message.CreatedAt,
message.SequenceId,
message.ForwardedFromId,
message.ForwardedFromId.HasValue && senders.TryGetValue(message.ForwardedFromId.Value, out var fwdUser) ? new MessageSenderDto(fwdUser.Id, fwdUser.Username, fwdUser.DisplayName, fwdUser.Avatar) : null,
storyMessage?.StoryId,
storyMessage?.StoryMediaUrl,
storyMessage?.StoryMediaType,
mediaMessage?.Media.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList() ?? new List<MediaDto>(),
senders.TryGetValue(message.SenderId, out var senderUser) ? new MessageSenderDto(senderUser.Id, senderUser.Username, senderUser.DisplayName, senderUser.Avatar) : null,
chat.Members.Where(m => m.LastReadSequenceId >= message.SequenceId && m.UserId != message.SenderId).Select(m => new ReadByDto(m.UserId)).ToList(),
reactionsWithUser,
null, // ForwardedFrom details not implemented here yet
(message as StoryMessage)?.StoryId,
(message as StoryMessage)?.StoryMediaUrl,
(message as StoryMessage)?.StoryMediaType,
(message as MediaMessage)?.Media.Select(m => new MediaDto(m.Id, m.Type, m.Url, m.Filename, m.Size, m.Duration)).ToList() ?? new List<MediaDto>(),
sender != null ? new MessageSenderDto(sender.Id, sender.Username, sender.DisplayName, sender.Avatar) : new MessageSenderDto(message.SenderId, "unknown", "Unknown", null),
new List<ReadByDto>(), // ReadBy not implemented in this detailed view yet
reactions?.Select(r => {
senders.TryGetValue(r.UserId, out var ru);
return new MessageReactionDto(r.Id, r.Emoji, r.UserId, ru != null ? new MessageSenderDto(ru.Id, ru.Username, ru.DisplayName, ru.Avatar) : null);
}).ToList() ?? new List<MessageReactionDto>(),
(message as CallMessage)?.CallType,
(message as CallMessage)?.CallStatus,
(message as CallMessage)?.Duration
(message as CallMessage)?.Duration,
(message as PollMessage)?.Options.Select(o => {
var pm = (PollMessage)message;
var voters = pm.IsAnonymous == false
? pm.Votes
.Where(v => v.OptionId == o.Id)
.Select(v => {
senders.TryGetValue(v.UserId, out var vu);
return vu != null
? new MessageSenderDto(vu.Id, vu.Username, vu.DisplayName, vu.Avatar)
: new MessageSenderDto(v.UserId, "unknown", "Unknown", null);
})
.ToList()
: null;
return new PollOptionDto(o.Id, o.Text, o.VoteCount, voters, pm.IsAnonymous == false ? pm.Votes.Where(v => v.OptionId == o.Id).Select(v => v.UserId).ToList() : null);
}).ToList(),
(message as PollMessage)?.IsMultipleChoice,
(message as PollMessage)?.IsAnonymous,
(message as PollMessage)?.IsClosed,
(message as PollMessage)?.Votes.Where(v => v.UserId == request.UserId).Select(v => v.OptionId).ToList()
));
}
return Result.Success(result);
}
}
@@ -6,9 +6,9 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
@@ -88,22 +88,17 @@ internal sealed class GetSharedMediaQueryHandler : IQueryHandler<GetSharedMediaQ
var filteredMedia = messageMedia.Where(media =>
{
var mType = media.Type?.ToLower() ?? "file";
var isGif = mType == "image" && media.Url != null && (media.Url.EndsWith(".mp4", StringComparison.OrdinalIgnoreCase) || media.Url.EndsWith(".gif", StringComparison.OrdinalIgnoreCase));
var filename = media.Filename?.ToLower() ?? "";
var url = media.Url?.ToLower() ?? "";
if (filterType == "gifs")
{
return isGif;
}
var isGif = mType == "gif" ||
(mType == "image" && (filename.EndsWith(".mp4") || filename.EndsWith(".gif") || url.EndsWith(".gif") || filename.Contains("gif"))) ||
(mType == "video" && (filename.Contains("animation") || filename.Contains("gif")));
if (filterType == "files")
{
return mType != "image" && mType != "video" && mType != "link";
}
if (filterType == "media")
{
return (mType == "image" || mType == "video") && !isGif;
}
if (filterType == "gifs") return isGif;
if (filterType == "media") return (mType == "image" || mType == "video") && !isGif;
if (filterType == "files") return (mType == "file" || mType == "audio") && !isGif && mType != "image" && mType != "video";
if (filterType == "links") return mType == "link";
return true;
}).ToList();
@@ -124,7 +119,7 @@ internal sealed class GetSharedMediaQueryHandler : IQueryHandler<GetSharedMediaQ
storyMessage?.StoryMediaType,
message.IsEdited,
message.Type,
filteredMedia.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size)).ToList()
filteredMedia.Select(media => new MediaDto(media.Id, media.Type, media.Url, media.Filename, media.Size, media.Duration)).ToList()
));
}
}
@@ -0,0 +1,58 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Conversations.Application.Messages.Pin;
public sealed record PinMessageCommand(Guid MessageId, Guid ChatId, Guid UserId) : ICommand<Guid>;
public sealed class PinMessageCommandHandler : ICommandHandler<PinMessageCommand, Guid>
{
private readonly IMessageRepository _messageRepository;
private readonly IChatRepository _chatRepository;
private readonly IChatsUnitOfWork _unitOfWork;
private readonly IMediator _mediator;
public PinMessageCommandHandler(
IMessageRepository messageRepository,
IChatRepository chatRepository,
IChatsUnitOfWork unitOfWork,
IMediator mediator)
{
_messageRepository = messageRepository;
_chatRepository = chatRepository;
_unitOfWork = unitOfWork;
_mediator = mediator;
}
public async Task<Result<Guid>> Handle(PinMessageCommand request, CancellationToken cancellationToken)
{
var chat = await _chatRepository.GetByIdAsync(request.ChatId, cancellationToken);
if (chat is null) return Result.Failure<Guid>(ChatErrors.ChatsNotFound);
// Security check
if (!chat.Members.Any(m => m.UserId == request.UserId))
return Result.Failure<Guid>(ChatErrors.ChatsForbidden);
var message = await _messageRepository.GetByIdAsync(request.MessageId, cancellationToken);
if (message is null) return Result.Failure<Guid>(ChatErrors.NotFound);
if (message.ChatId != request.ChatId)
return Result.Failure<Guid>(ChatErrors.NotFound);
message.AddState(MessageState.IsPinned);
await _messageRepository.UpdateAsync(message, cancellationToken);
await _unitOfWork.SaveChangesAsync(cancellationToken);
// Notify chat about pinned message change
await _mediator.Publish(new MessagePinnedDomainEvent(message.Id, message.ChatId, message.SenderId, message.Content), cancellationToken);
return Result.Success(message.Id);
}
}
public record MessagePinnedDomainEvent(Guid MessageId, Guid ChatId, Guid SenderId, string? Content) : INotification;
@@ -1,7 +1,7 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.SignalR;
using Knot.Shared.Kernel;
using MediatR;
@@ -1,7 +1,7 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.SignalR;
using Knot.Shared.Kernel;
using MediatR;
@@ -1,7 +1,7 @@
using MediatR;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
namespace Knot.Modules.Conversations.Application.Messages.Read;
@@ -6,7 +6,7 @@ using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
@@ -1,8 +1,8 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Settings.Application.Abstractions;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Application.Messages.Send;
@@ -135,8 +135,9 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
else if (request.Type == "poll")
{
if (!_messagesSettings.Current.AllowPolls) return Result.Failure<Guid>(ChatErrors.PollsDisabled);
if (chat.Type != ChatType.Group) return Result.Failure<Guid>(new Error("Poll.InvalidChat", "Polls are only allowed in groups."));
message = new PollMessage(
message = PollMessage.Create(
Guid.NewGuid(),
request.ChatId,
request.SenderId,
@@ -146,9 +147,7 @@ public sealed class SendMessageCommandHandler : ICommandHandler<SendMessageComma
request.PollAllowMultipleAnswers ?? false,
request.PollExpiresAt,
request.ReplyToId,
request.ForwardedFromId,
DateTime.UtcNow,
false);
request.ForwardedFromId);
}
else if (request.Type == "call")
{
@@ -0,0 +1,58 @@
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Conversations.Application.Messages.Unpin;
public sealed record UnpinMessageCommand(Guid MessageId, Guid ChatId, Guid UserId) : ICommand<Guid>;
public sealed class UnpinMessageCommandHandler : ICommandHandler<UnpinMessageCommand, Guid>
{
private readonly IMessageRepository _messageRepository;
private readonly IChatRepository _chatRepository;
private readonly IChatsUnitOfWork _unitOfWork;
private readonly IMediator _mediator;
public UnpinMessageCommandHandler(
IMessageRepository messageRepository,
IChatRepository chatRepository,
IChatsUnitOfWork unitOfWork,
IMediator mediator)
{
_messageRepository = messageRepository;
_chatRepository = chatRepository;
_unitOfWork = unitOfWork;
_mediator = mediator;
}
public async Task<Result<Guid>> Handle(UnpinMessageCommand request, CancellationToken cancellationToken)
{
var chat = await _chatRepository.GetByIdAsync(request.ChatId, cancellationToken);
if (chat is null) return Result.Failure<Guid>(ChatErrors.ChatsNotFound);
// Security check
if (!chat.Members.Any(m => m.UserId == request.UserId))
return Result.Failure<Guid>(ChatErrors.ChatsForbidden);
var message = await _messageRepository.GetByIdAsync(request.MessageId, cancellationToken);
if (message is null) return Result.Failure<Guid>(ChatErrors.NotFound);
if (message.ChatId != request.ChatId)
return Result.Failure<Guid>(ChatErrors.NotFound);
message.RemoveState(MessageState.IsPinned);
await _messageRepository.UpdateAsync(message, cancellationToken);
await _unitOfWork.SaveChangesAsync(cancellationToken);
// Notify chat about unpinned message change
await _mediator.Publish(new MessageUnpinnedDomainEvent(message.Id, message.ChatId), cancellationToken);
return Result.Success(message.Id);
}
}
public record MessageUnpinnedDomainEvent(Guid MessageId, Guid ChatId) : INotification;
@@ -1,4 +1,4 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using System;
using System.IO;
using System.Threading;
@@ -0,0 +1,110 @@
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Shared.Kernel;
using MediatR;
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using System.Collections.Generic;
namespace Knot.Modules.Conversations.Application.Messages.Vote;
public sealed record VotePollCommand(
Guid MessageId,
Guid ChatId,
Guid UserId,
Guid OptionId) : ICommand;
public sealed class VotePollCommandHandler : ICommandHandler<VotePollCommand>
{
private readonly IMessageRepository _messageRepository;
private readonly IChatRepository _chatRepository;
private readonly IChatsUnitOfWork _unitOfWork;
private readonly IMessageNotifier _notifier;
private readonly IUserDisplayNameProvider _userProvider;
public VotePollCommandHandler(
IMessageRepository messageRepository,
IChatRepository chatRepository,
IChatsUnitOfWork unitOfWork,
IMessageNotifier notifier,
IUserDisplayNameProvider userProvider)
{
_messageRepository = messageRepository;
_chatRepository = chatRepository;
_unitOfWork = unitOfWork;
_notifier = notifier;
_userProvider = userProvider;
}
public async Task<Result> Handle(VotePollCommand request, CancellationToken cancellationToken)
{
var message = await _messageRepository.GetByIdAsync(request.MessageId, cancellationToken);
if (message is not PollMessage poll) return Result.Failure(new Error("Poll.NotFound", "Poll not found"));
if (poll.IsClosed) return Result.Failure(new Error("Poll.Closed", "This poll is closed."));
var chat = await _chatRepository.GetByIdAsync(request.ChatId, cancellationToken);
if (chat == null || !chat.Members.Any(m => m.UserId == request.UserId)) return Result.Failure(ChatErrors.ChatsForbidden);
var targetOption = poll.Options.FirstOrDefault(o => o.Id == request.OptionId);
if (targetOption == null) return Result.Failure(new Error("Poll.InvalidOption", "Invalid option ID."));
// Prevent duplicate or changed votes
var existingVote = poll.Votes.FirstOrDefault(v => v.UserId == request.UserId && v.OptionId == request.OptionId);
if (existingVote != null) return Result.Failure(new Error("Poll.AlreadyVoted", "You have already voted for this option."));
if (!poll.IsMultipleChoice)
{
var hasVotedInThisPoll = poll.Votes.Any(v => v.UserId == request.UserId);
if (hasVotedInThisPoll) return Result.Failure(new Error("Poll.AlreadyVoted", "You have already voted in this poll."));
}
poll.Votes.Add(new PollVote { UserId = request.UserId, OptionId = request.OptionId, VotedAt = DateTime.UtcNow });
targetOption.VoteCount++;
await _messageRepository.UpdateAsync(poll, cancellationToken);
await _unitOfWork.SaveChangesAsync(cancellationToken);
// Notify updated poll
var voterIds = poll.Votes.Select(v => v.UserId).Distinct().ToList();
var votersInfo = poll.IsAnonymous == false
? await _userProvider.GetUsersInfoAsync(voterIds, cancellationToken)
: new Dictionary<Guid, UserInfo>();
await _notifier.NotifyMessageUpdateAsync(poll.ChatId, "poll_updated", new
{
id = poll.Id,
chatId = poll.ChatId,
senderId = poll.SenderId,
createdAt = poll.CreatedAt,
type = "poll",
content = poll.Content,
pollOptions = poll.Options.Select(o => new {
id = o.Id,
text = o.Text,
voteCount = o.VoteCount,
voters = poll.IsAnonymous == false
? poll.Votes.Where(v => v.OptionId == o.Id)
.Select(v => {
votersInfo.TryGetValue(v.UserId, out var vu);
return vu != null
? new { id = vu.Id, username = vu.Username, displayName = vu.DisplayName, avatar = vu.Avatar }
: new { id = v.UserId, username = "unknown", displayName = "Unknown", avatar = (string?)null };
}).ToList()
: null,
voterIds = poll.IsAnonymous == false
? poll.Votes.Where(v => v.OptionId == o.Id).Select(v => v.UserId).ToList()
: null
}).ToList(),
pollIsMultipleChoice = poll.IsMultipleChoice,
pollIsClosed = poll.IsClosed,
pollIsAnonymous = poll.IsAnonymous
}, cancellationToken);
return Result.Success();
}
}
@@ -2,8 +2,8 @@ using System.Text.RegularExpressions;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Contracts.Auth.Domain;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Storage;
using MediatR;
@@ -54,9 +54,11 @@ internal sealed class DeleteUserCommandHandler : ICommandHandler<DeleteUserComma
if (msg is MediaMessage mediaMsg)
{
foreach (var media in mediaMsg.Media)
{
if (!string.IsNullOrEmpty(media.Url))
{
bool isUsedElsewhere = allMessages.Any(m => m.Id != msg.Id &&
m is MediaMessage mm && mm.Media.Any(ame => ame.Url == media.Url));
m is MediaMessage mm && mm.Media.Any(ame => !string.IsNullOrEmpty(ame.Url) && ame.Url == media.Url));
if (!isUsedElsewhere)
{
@@ -69,6 +71,7 @@ internal sealed class DeleteUserCommandHandler : ICommandHandler<DeleteUserComma
}
}
}
}
await _messageRepository.DeleteUserMessagesAsync(request.UserId, cancellationToken);
@@ -1,16 +1,14 @@
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Infrastructure.Persistence;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.Persistence;
using Knot.Modules.Conversations.Infrastructure.Persistence.Mongo;
using Knot.Modules.Conversations.Infrastructure.Services;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.Configuration;
using Microsoft.Extensions.DependencyInjection;
using ConversationsAbstractions = Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Messaging.Application.Abstractions;
namespace Knot.Modules.Conversations;
@@ -29,21 +27,22 @@ public static class DependencyInjection
var mongoConnectionString = configuration.GetConnectionString("MongoConnection") ?? "mongodb://localhost:27017";
services.AddScoped<ConversationsAbstractions.IChatsUnitOfWork>(sp => sp.GetRequiredService<ChatsDbContext>());
services.AddScoped<IChatsUnitOfWork>(sp => sp.GetRequiredService<ChatsDbContext>());
services.AddScoped<Knot.Contracts.Conversations.Infrastructure.Persistence.IChatsDbContext>(sp => sp.GetRequiredService<ChatsDbContext>());
services.AddScoped<IChatRepository, ChatRepository>();
services.AddScoped<IFolderRepository, FolderRepository>();
services.AddScoped<IUserChatSettingsRepository, UserChatSettingsRepository>();
services.AddScoped<IUserFolderSettingsRepository, UserFolderSettingsRepository>();
services.AddScoped<Knot.Contracts.Conversations.Domain.IUserFolderSettingsRepository, UserFolderSettingsRepository>();
services.AddMediatR(config =>
config.RegisterServicesFromAssembly(typeof(DependencyInjection).Assembly));
services.AddScoped<Knot.Contracts.Messaging.Application.Abstractions.IChatAccessProvider, Knot.Modules.Conversations.Infrastructure.Services.ChatAccessProvider>();
services.AddScoped<Knot.Contracts.Messaging.Application.Abstractions.IChatAccessProvider, ChatAccessProvider>();
services.AddScoped<Knot.Contracts.Messaging.Application.Abstractions.IMessageNotifier, Knot.Modules.Conversations.Infrastructure.SignalR.MessageNotifier>();
services.AddScoped<ConversationsAbstractions.IUserStatusService, Knot.Modules.Conversations.Infrastructure.Services.UserStatusService>();
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserStatusService, Knot.Modules.Conversations.Infrastructure.Services.UserStatusService>();
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserDeleterService, Knot.Modules.Conversations.Infrastructure.Services.UserDeleterService>();
services.AddScoped<Knot.Contracts.Conversations.Application.Abstractions.IUserStatusService, UserStatusService>();
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserStatusService, UserStatusService>();
services.AddScoped<Knot.Contracts.Conversations.Abstractions.IUserDeleterService, UserDeleterService>();
return services;
}
}
@@ -1,157 +0,0 @@
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Domain;
public sealed record ChatCreatedDomainEvent(Chat Chat) : IDomainEvent;
public sealed record ChatMemberAddedDomainEvent(Guid ChatId, Guid UserId) : IDomainEvent;
/// <summary>
/// Тип чата: личный или групповой.
/// </summary>
public enum ChatType
{
Personal,
Group,
Favorites
}
/// <summary>
/// Роль участника в чате.
/// </summary>
public static class ChatRole
{
public const string Owner = "owner";
public const string Admin = "admin";
public const string Member = "member";
}
/// <summary>
/// Сущность чата (Агрегат).
/// </summary>
public sealed class Chat : AggregateRoot<Guid>
{
public ChatType Type { get; private set; }
public string? Name { get; private set; }
public string? Description { get; private set; }
public string? Avatar { get; private set; }
public DateTime CreatedAt { get; private set; }
public long LastMessageSequenceId { get; private set; }
private readonly List<ChatMember> _members = new();
public IReadOnlyCollection<ChatMember> Members => _members.AsReadOnly();
private Chat(Guid id, ChatType type, string? name, string? avatar, string? description = null) : base(id)
{
Type = type;
Name = name;
Avatar = avatar;
Description = description;
CreatedAt = DateTime.UtcNow;
}
/// <summary>
/// Создает личный чат между двумя пользователями.
/// </summary>
public static Chat CreatePersonal()
{
var chat = new Chat(Guid.NewGuid(), ChatType.Personal, null, null);
chat.RaiseDomainEvent(new ChatCreatedDomainEvent(chat));
return chat;
}
/// <summary>
/// Создает групповой чат.
/// </summary>
public static Chat CreateGroup(string name, string? avatar = null)
{
var chat = new Chat(Guid.NewGuid(), ChatType.Group, name, avatar);
chat.RaiseDomainEvent(new ChatCreatedDomainEvent(chat));
return chat;
}
/// <summary>
/// Фабричный метод для создания чата.
/// </summary>
public static Chat Create(string? name, ChatType type, string? avatar = null, string? description = null)
{
var chat = new Chat(Guid.NewGuid(), type, name, avatar, description);
chat.RaiseDomainEvent(new ChatCreatedDomainEvent(chat));
return chat;
}
public void AddMember(Guid userId, string role = "member")
{
if (_members.Any(m => m.UserId == userId))
{
return;
}
_members.Add(new ChatMember(Id, userId, role));
RaiseDomainEvent(new ChatMemberAddedDomainEvent(Id, userId));
}
public void RemoveMember(Guid userId)
{
var member = _members.FirstOrDefault(m => m.UserId == userId);
if (member != null)
{
_members.Remove(member);
}
}
public void UpdateName(string name) => Name = name;
public void UpdateDescription(string? description) => Description = description;
public void UpdateAvatar(string? avatarUrl) => Avatar = avatarUrl;
public long IncrementSequenceId()
{
return ++LastMessageSequenceId;
}
}
/// <summary>
/// Участник чата.
/// </summary>
public sealed class ChatMember : Entity<Guid>
{
public Guid ChatId { get; private set; }
public Guid UserId { get; private set; }
public string Role { get; private set; }
public DateTime JoinedAt { get; private set; }
public bool IsPinned { get; private set; }
public bool IsMuted { get; private set; }
public Guid? LastReadMessageId { get; private set; }
public long LastReadSequenceId { get; private set; }
public Guid? LastDeliveredMessageId { get; private set; }
// For EF Core
private ChatMember() : base(Guid.Empty) { Role = "member"; }
internal ChatMember(Guid chatId, Guid userId, string role) : base(Guid.NewGuid())
{
ChatId = chatId;
UserId = userId;
Role = role;
JoinedAt = DateTime.UtcNow;
}
public void TogglePin() => IsPinned = !IsPinned;
public void UpdateReadCursor(Guid messageId, long sequenceId)
{
if (sequenceId > LastReadSequenceId)
{
LastReadMessageId = messageId;
LastReadSequenceId = sequenceId;
}
}
public void UpdateDeliveredCursor(Guid messageId)
{
LastDeliveredMessageId = messageId;
}
}
@@ -1,10 +0,0 @@
namespace Knot.Modules.Conversations.Domain;
public static class ChatConstants
{
public const int DefaultMessageQueryLimit = 100;
public const int MaxSharedMediaQueryLimit = 300;
public const int SearchMessagesLimit = 50;
public const int MaxFileUploadSizeMb = 50;
}
@@ -1,25 +0,0 @@
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Domain;
public static class ChatErrors
{
public static readonly Error FileEmpty = new Error("File.Empty", "No file uploaded");
public static readonly Error FileInvalidExtension = new Error("File.InvalidExtension", "Must be a ZIP archive");
public static readonly Error ImportExpired = new Error("Import.Expired", "Session not found or expired");
public static readonly Error ImportMissing = new Error("Import.Missing", "ZIP file lost");
public static readonly Error ChatNotFound = new Error("Chat.NotFound", "Chat not found or access denied");
public static readonly Error NotFound = new Error("Chat.NotFound", "Chat not found"); // Alias
public static readonly Error NotMember = new Error("Chat.NotMember", "You are not a member of this chat");
public static readonly Error ChatsForbidden = new Error("Chats.Forbidden", "Вы не являетесь участником этого чата.");
public static readonly Error MessagesNotFound = new Error("Messages.NotFound", "Message not found.");
public static readonly Error ChatsNotFound = new Error("Chats.NotFound", "Чат не найден.");
public static readonly Error Unauthorized = new Error("Chats.Unauthorized", "Access denied");
public static readonly Error FoldersDisabled = new Error("Folders.Disabled", "Folders feature is disabled by the administrator.");
public static readonly Error PollsDisabled = new Error("Polls.Disabled", "Polls are disabled by the administrator.");
public static readonly Error MediaDisabled = new Error("Media.Disabled", "Media messages are disabled by the administrator.");
public static Error ImportCreateChatFailed(string msg) => new Error("Import.CreateChatFailed", msg);
public static Error FileTooLarge(int maxMb) => new Error("File.TooLarge", $"File exceeds the maximum allowed size of {maxMb}MB.");
}
@@ -1,108 +0,0 @@
using System;
using System.Collections.Generic;
using Knot.Shared.Kernel;
namespace Knot.Modules.Conversations.Domain;
/// <summary>
/// Сущность папки для группировки чатов.
/// </summary>
public sealed class Folder : AggregateRoot<Guid>
{
public string Name { get; private set; }
public string? Icon { get; private set; } // URL из хранилища
public bool IsDefault { get; private set; }
public FolderType Type { get; private set; }
public Folder(Guid id, string name, string? icon = null, bool isDefault = false, FolderType type = FolderType.Custom)
: base(id)
{
Name = name;
Icon = icon;
IsDefault = isDefault;
Type = type;
}
public void Update(string name, string? icon)
{
if (IsDefault) throw new InvalidOperationException("Cannot rename default folders.");
Name = name;
Icon = icon;
}
}
public enum FolderType
{
All, // Все чаты
New, // Новые (с непрочитанными)
Muted, // Без звука
Custom // Пользовательская
}
/// <summary>
/// Настройки конкретного чата для конкретного пользователя.
/// Хранятся в PostgreSQL (связь User <-> Chat).
/// </summary>
public sealed class UserChatSettings : Entity<Guid>
{
public Guid UserId { get; private set; }
public Guid ChatId { get; private set; }
// Список папок, в которые входит чат для этого пользователя
private readonly List<Guid> _folderIds = new();
public IReadOnlyCollection<Guid> FolderIds => _folderIds.AsReadOnly();
public bool IsMuted { get; private set; }
private UserChatSettings() : base(Guid.NewGuid()) { }
public UserChatSettings(Guid userId, Guid chatId) : base(Guid.NewGuid())
{
UserId = userId;
ChatId = chatId;
}
public static UserChatSettings Create(Guid userId, Guid chatId) => new(userId, chatId);
public void AddToFolder(Guid folderId)
{
if (!_folderIds.Contains(folderId)) _folderIds.Add(folderId);
}
public void RemoveFromFolder(Guid folderId)
{
_folderIds.Remove(folderId);
}
public void SetMute(bool isMuted) => IsMuted = isMuted;
}
/// <summary>
/// Глобальные настройки папок пользователя (скрытие дефолтных и т.д.).
/// Будет храниться в MongoDB.
/// </summary>
public sealed class UserFolderSettings : AggregateRoot<Guid>
{
public Guid UserId { get; private set; }
// Список ID папок, которые пользователь скрыл (только для дефолтных)
public List<Guid> HiddenDefaultFolderIds { get; private set; } = new();
// Список пользовательских папок (Guid созданных Folder)
public List<Guid> CustomFolderIds { get; private set; } = new();
public UserFolderSettings(Guid userId) : base(Guid.NewGuid())
{
UserId = userId;
}
public void HideFolder(Guid folderId)
{
if (!HiddenDefaultFolderIds.Contains(folderId)) HiddenDefaultFolderIds.Add(folderId);
}
public void ShowFolder(Guid folderId)
{
HiddenDefaultFolderIds.Remove(folderId);
}
}
@@ -1,39 +0,0 @@
using Knot.Modules.Conversations.Domain;
namespace Knot.Modules.Conversations.Domain;
public interface IChatRepository
{
void Add(Chat chat);
void Update(Chat chat);
void Remove(Chat chat);
Task<Chat?> GetByIdAsync(Guid id, CancellationToken cancellationToken);
Task<Chat?> GetFavoritesAsync(Guid userId, CancellationToken cancellationToken);
Task<List<Chat>> GetUserChatsAsync(Guid userId, CancellationToken cancellationToken);
}
public interface IFolderRepository
{
void Add(Folder folder);
void Update(Folder folder);
void Remove(Folder folder);
Task<Folder?> GetByIdAsync(Guid id, CancellationToken cancellationToken);
Task<List<Folder>> GetUserFoldersAsync(Guid userId, CancellationToken cancellationToken);
}
public interface IUserChatSettingsRepository
{
void Add(UserChatSettings settings);
void Update(UserChatSettings settings);
void Remove(UserChatSettings settings);
void RemoveRange(IEnumerable<UserChatSettings> settings);
Task<UserChatSettings?> GetAsync(Guid userId, Guid chatId, CancellationToken cancellationToken);
Task<List<UserChatSettings>> GetByUserIdAsync(Guid userId, CancellationToken cancellationToken);
}
public interface IUserFolderSettingsRepository
{
Task<UserFolderSettings?> GetByUserIdAsync(Guid userId, CancellationToken cancellationToken);
Task UpdateAsync(UserFolderSettings settings, CancellationToken cancellationToken);
Task RemoveByUserIdAsync(Guid userId, CancellationToken cancellationToken);
}
@@ -1,5 +1,5 @@
using Knot.Shared.Kernel;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Modules.Conversations.Infrastructure.SignalR;
using MediatR;
using Microsoft.AspNetCore.SignalR;
@@ -1,5 +1,5 @@
using Knot.Contracts.Conversations.Domain;
using Microsoft.EntityFrameworkCore;
using Knot.Modules.Conversations.Domain;
namespace Knot.Modules.Conversations.Infrastructure.Persistence;
@@ -4,19 +4,18 @@ using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Domain;
using Knot.Contracts.Conversations.Infrastructure.Persistence;
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Domain;
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Security;
using MediatR;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using DomainChat = Knot.Modules.Conversations.Domain.Chat;
using DomainChat = Knot.Contracts.Conversations.Domain.Chat;
namespace Knot.Modules.Conversations.Infrastructure.Persistence;
public sealed class ChatsDbContext : DbContext, Knot.Modules.Conversations.Application.Abstractions.IChatsUnitOfWork, Knot.Contracts.Conversations.Infrastructure.Persistence.IChatsDbContext
public sealed class ChatsDbContext : DbContext, Knot.Contracts.Conversations.Application.Abstractions.IChatsUnitOfWork, Knot.Contracts.Conversations.Infrastructure.Persistence.IChatsDbContext
{
private readonly IMediator? _mediator;
private readonly IEncryptionService? _encryptionService;
@@ -54,6 +53,8 @@ public sealed class ChatsDbContext : DbContext, Knot.Modules.Conversations.Appli
{
builder.ToTable("Chats");
builder.HasKey(c => c.Id);
builder.Property(c => c.IsImporting);
builder.Property(c => c.ImportJobId);
builder.Property(c => c.Type).HasConversion<string>();
builder.OwnsMany(c => c.Members, mb =>
@@ -1,4 +1,4 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using MongoDB.Bson.Serialization;
namespace Knot.Modules.Conversations.Infrastructure.Persistence.Mongo;
@@ -1,4 +1,4 @@
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using MongoDB.Driver;
namespace Knot.Modules.Conversations.Infrastructure.Persistence.Mongo;
@@ -1,9 +1,11 @@
using Knot.Modules.Conversations.Application.Abstractions;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Modules.Conversations.Infrastructure.SignalR;
namespace Knot.Modules.Conversations.Infrastructure.Services;
public sealed class UserStatusService : IUserStatusService, Knot.Contracts.Conversations.Abstractions.IUserStatusService
public sealed class UserStatusService :
Knot.Contracts.Conversations.Application.Abstractions.IUserStatusService,
Knot.Contracts.Conversations.Abstractions.IUserStatusService
{
public bool IsUserOnline(string userId)
{
@@ -8,11 +8,20 @@ using Knot.Modules.Conversations.Application.Messages.Send;
using Knot.Modules.Conversations.Application.Messages.Read;
using Knot.Modules.Conversations.Application.Messages.Delete;
using Knot.Modules.Conversations.Application.Messages.React;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using Microsoft.Extensions.Caching.Memory;
using Knot.Contracts.Auth.Domain;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Modules.Conversations.Application.Messages.Pin;
using Knot.Modules.Conversations.Application.Messages.Unpin;
using Knot.Modules.Conversations.Application.Messages.Vote;
using Knot.Modules.Conversations.Application.Messages.Edit;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Contracts.Messaging.Application.Abstractions;
using Knot.Contracts.Messaging.Domain;
using Knot.Contracts.Conversations.Application.Abstractions;
using Knot.Contracts.Profiles.Domain;
namespace Knot.Modules.Conversations.Infrastructure.SignalR;
@@ -39,17 +48,32 @@ public sealed class ChatHub : Hub
private readonly IUserContext _userContext;
private readonly IChatRepository _chatRepository;
private readonly IUserRepository _userRepository;
private readonly IMessageRepository _messageRepository;
private readonly ILogger<ChatHub> _logger;
private readonly IMemoryCache _cache;
private readonly IUserDisplayNameProvider _userProvider;
private readonly IProfileRepository _profileRepository;
public ChatHub(ISender sender, IUserContext userContext, IChatRepository chatRepository, IUserRepository userRepository, ILogger<ChatHub> logger, IMemoryCache cache)
public ChatHub(
ISender sender,
IUserContext userContext,
IChatRepository chatRepository,
IUserRepository userRepository,
IMessageRepository messageRepository,
ILogger<ChatHub> logger,
IMemoryCache cache,
IUserDisplayNameProvider userProvider,
IProfileRepository profileRepository)
{
_sender = sender;
_userContext = userContext;
_chatRepository = chatRepository;
_userRepository = userRepository;
_messageRepository = messageRepository;
_logger = logger;
_cache = cache;
_userProvider = userProvider;
_profileRepository = profileRepository;
}
public override async Task OnConnectedAsync()
@@ -75,7 +99,15 @@ public sealed class ChatHub : Hub
userId, Context.ConnectionId, userChats.Count);
await Clients.Others.SendAsync("user_online", new { userId });
var isInvisible = await IsUserInvisibleAsync(_userContext.UserId, Context.ConnectionAborted);
if (!isInvisible)
{
await BroadcastPresenceToVisibleUsersAsync(
"user_online",
new { userId },
_userContext.UserId,
Context.ConnectionAborted);
}
}
await base.OnConnectedAsync();
}
@@ -91,7 +123,22 @@ public sealed class ChatHub : Hub
if (set.Count == 0)
{
_userConnections.TryRemove(userId, out _);
await Clients.Others.SendAsync("user_offline", new { userId, lastSeen = DateTime.UtcNow });
try
{
var isInvisible = await IsUserInvisibleAsync(_userContext.UserId, CancellationToken.None);
if (!isInvisible)
{
await BroadcastPresenceToVisibleUsersAsync(
"user_offline",
new { userId, lastSeen = DateTime.UtcNow },
_userContext.UserId,
CancellationToken.None);
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Presence broadcast on disconnect failed for {UserId}", userId);
}
}
}
_cache.Set("Global_OnlineUsersCount", _userConnections.Count);
@@ -100,6 +147,56 @@ public sealed class ChatHub : Hub
await base.OnDisconnectedAsync(exception);
}
private async Task<bool> IsUserInvisibleAsync(Guid userId, CancellationToken ct)
{
var profile = await _profileRepository.GetAsync(userId, ct);
return profile?.IsInvisible ?? false;
}
private async Task BroadcastPresenceToVisibleUsersAsync(
string eventName,
object payload,
Guid sourceUserId,
CancellationToken cancellationToken)
{
var recipients = _userConnections.Keys
.Where(id => id != sourceUserId.ToString())
.Select(id => Guid.TryParse(id, out var parsed) ? parsed : Guid.Empty)
.Where(id => id != Guid.Empty)
.Distinct()
.ToList();
if (recipients.Count == 0)
return;
var recipientProfiles = await _profileRepository.GetAsync(recipients, cancellationToken);
var invisibleRecipientIds = recipientProfiles
.Where(p => p.IsInvisible)
.Select(p => p.UserId)
.ToHashSet();
foreach (var kvp in _userConnections)
{
if (!Guid.TryParse(kvp.Key, out var recipientId))
continue;
if (recipientId == sourceUserId)
continue;
if (invisibleRecipientIds.Contains(recipientId))
continue;
string[] connectionIds;
lock (kvp.Value)
{
connectionIds = kvp.Value.ToArray();
}
foreach (var connectionId in connectionIds)
{
await Clients.Client(connectionId).SendAsync(eventName, payload);
}
}
}
// ────────────────────────────────────────────────────────────────
// Chat methods
// ────────────────────────────────────────────────────────────────
@@ -112,14 +209,18 @@ public sealed class ChatHub : Hub
new AttachmentRequest(a.Type, a.Url, a.FileName, a.FileSize)).ToList();
var command = new SendMessageCommand(
request.ChatId,
_userContext.UserId,
request.Content,
request.Type,
attachments,
request.ReplyToId,
request.Quote,
request.ForwardedFromId);
ChatId: request.ChatId,
SenderId: _userContext.UserId,
Content: request.Content,
Type: request.Type,
Attachments: attachments,
ReplyToId: request.ReplyToId,
Quote: request.Quote,
ForwardedFromId: request.ForwardedFromId,
PollOptions: request.PollOptions,
PollIsAnonymous: request.PollIsAnonymous,
PollAllowMultipleAnswers: request.PollAllowMultipleAnswers
);
await _sender.Send(command);
}
@@ -227,6 +328,62 @@ public sealed class ChatHub : Hub
_logger.LogInformation("RemoveReaction completed successfully");
}
[HubMethodName("pin_message")]
public async Task PinMessage(PinMessageRequest request)
{
var command = new PinMessageCommand(request.MessageId, request.ChatId, _userContext.UserId);
var result = await _sender.Send(command);
var message = await _messageRepository.GetByIdAsync(request.MessageId, Context.ConnectionAborted);
if (message != null)
{
var senderInfo = await _userProvider.GetUsersInfoAsync(new[] { message.SenderId });
var dto = MessageMapper.MapToDto(message, senderInfo, Enumerable.Empty<MessageReaction>(), Enumerable.Empty<Guid>());
await Clients.Group(request.ChatId.ToString()).SendAsync("message_pinned", new
{
chatId = request.ChatId,
message = dto,
userId = _userContext.UserId
});
}
}
[HubMethodName("unpin_message")]
public async Task UnpinMessage(PinMessageRequest request)
{
var command = new UnpinMessageCommand(request.MessageId, request.ChatId, _userContext.UserId);
var result = await _sender.Send(command);
await Clients.Group(request.ChatId.ToString()).SendAsync("message_unpinned", new
{
chatId = request.ChatId,
messageId = request.MessageId,
userId = _userContext.UserId
});
}
[HubMethodName("edit_message")]
public async Task EditMessage(EditMessageHubRequest request)
{
var command = new EditMessageCommand(request.MessageId, request.ChatId, _userContext.UserId, request.Content);
var result = await _sender.Send(command);
if (result.IsFailure)
{
throw new HubException(result.Error.Description);
}
}
[HubMethodName("vote_poll")]
public async Task VotePoll(VotePollRequest request)
{
var command = new VotePollCommand(request.MessageId, request.ChatId, _userContext.UserId, request.OptionId);
var result = await _sender.Send(command);
if (result.IsFailure)
{
throw new HubException(result.Error.Description);
}
}
// ────────────────────────────────────────────────────────────────
// Friend signals (Proxy methods for real-time notification)
// ────────────────────────────────────────────────────────────────
@@ -648,7 +805,7 @@ public sealed class ChatHub : Hub
await Clients.Group(request.ChatId).SendAsync("group_call_status_updated", new
{
chatId = request.ChatId,
userId = Context.UserIdentifier,
userId = _userContext.UserId.ToString(),
isMuted = request.IsMuted,
isVideoOff = request.IsVideoOff
});
@@ -675,7 +832,7 @@ public sealed class ChatHub : Hub
await Clients.Group(chatId).SendAsync("group_call_status_updated", new
{
chatId = chatId,
userId = Context.UserIdentifier,
userId = _userContext.UserId.ToString(),
isMuted = isMuted,
isVideoOff = isVideoOff
});
@@ -759,7 +916,10 @@ public sealed class ChatHub : Hub
List<AttachmentHubRequest>? Attachments = null,
Guid? ReplyToId = null,
string? Quote = null,
Guid? ForwardedFromId = null);
Guid? ForwardedFromId = null,
List<string>? PollOptions = null,
bool? PollIsAnonymous = null,
bool? PollAllowMultipleAnswers = null);
public record ReadMessagesRequest(Guid ChatId, Guid LastReadMessageId, long LastReadSequenceId);
public record CallOfferRequest(string TargetUserId, object Offer, string CallType, string? ChatId);
public record CallAnswerRequest(string TargetUserId, object Answer);
@@ -771,6 +931,7 @@ public sealed class ChatHub : Hub
public record CallStatusRequest(string TargetUserId, bool IsMuted, bool IsVideoOff);
public record AddReactionRequest(Guid MessageId, Guid ChatId, string Emoji);
public record RemoveReactionRequest(Guid MessageId, Guid ChatId, string Emoji);
public record PinMessageRequest(Guid MessageId, Guid ChatId);
public record DeleteMessagesHubRequest(Guid ChatId, List<string> MessageIds, bool DeleteForAll);
public record GroupCallJoinRequest(string ChatId, string CallType);
public record ParticipantInfo(string Id, string Username, string DisplayName, string? Avatar, bool IsSharingScreen = false, bool IsMuted = false, bool IsVideoOff = false);
@@ -783,6 +944,8 @@ public sealed class ChatHub : Hub
public record GroupRenegotiateRequest(string ChatId, string TargetUserId, object Offer);
public record GroupRenegotiateAnswerRequest(string ChatId, string TargetUserId, object Answer);
public record FriendSignalRequest(string FriendId);
public record VotePollRequest(Guid MessageId, Guid ChatId, Guid OptionId);
public record EditMessageHubRequest(Guid MessageId, Guid ChatId, string Content);
public class CallSession
{
@@ -5,4 +5,23 @@ using System.Threading;
using System.Threading.Tasks;
using Knot.Contracts.Messaging.Application.Abstractions;
using Microsoft.AspNetCore.SignalR;
public class MessageNotifier : IMessageNotifier { private readonly IHubContext<ChatHub> _hubContext; public MessageNotifier(IHubContext<ChatHub> hubContext) { _hubContext = hubContext; } public Task NotifyNewMessageAsync(Guid chatId, object messagePayload, CancellationToken cancellationToken) { return _hubContext.Clients.Group(chatId.ToString()).SendAsync("new_message", messagePayload, cancellationToken); } }
public class MessageNotifier : IMessageNotifier
{
private readonly IHubContext<ChatHub> _hubContext;
public MessageNotifier(IHubContext<ChatHub> hubContext)
{
_hubContext = hubContext;
}
public Task NotifyNewMessageAsync(Guid chatId, object messagePayload, CancellationToken cancellationToken)
{
return _hubContext.Clients.Group(chatId.ToString()).SendAsync("new_message", messagePayload, cancellationToken);
}
public Task NotifyMessageUpdateAsync(Guid chatId, string updateType, object updatePayload, CancellationToken cancellationToken)
{
return _hubContext.Clients.Group(chatId.ToString()).SendAsync(updateType, updatePayload, cancellationToken);
}
}
@@ -12,6 +12,7 @@
<ProjectReference Include="..\..\Contracts\Messaging\Knot.Contracts.Messaging.csproj" />
<ProjectReference Include="..\..\Contracts\Conversations\Knot.Contracts.Conversations.csproj" />
<ProjectReference Include="..\..\Contracts\Auth\Knot.Contracts.Auth.csproj" />
<ProjectReference Include="..\..\Contracts\Profiles\Knot.Contracts.Profiles.csproj" />
</ItemGroup>
<ItemGroup>
@@ -26,7 +26,7 @@ namespace Knot.Modules.Conversations.Migrations
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Modules.Conversations.Domain.Chat", b =>
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
@@ -56,9 +56,9 @@ namespace Knot.Modules.Conversations.Migrations
b.ToTable("Chats", "chats");
});
modelBuilder.Entity("Knot.Modules.Conversations.Domain.Chat", b =>
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.OwnsMany("Knot.Modules.Conversations.Domain.ChatMember", "Members", b1 =>
b.OwnsMany("Knot.Contracts.Conversations.Domain.ChatMember", "Members", b1 =>
{
b1.Property<Guid>("Id")
.ValueGeneratedOnAdd()
@@ -0,0 +1,166 @@
// <auto-generated />
using System;
using Knot.Modules.Conversations.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
#nullable disable
namespace Knot.Modules.Conversations.Migrations
{
[DbContext(typeof(ChatsDbContext))]
[Migration("20260406122703_AddIsImportingToChat")]
partial class AddIsImportingToChat
{
/// <inheritdoc />
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasDefaultSchema("chats")
.HasAnnotation("ProductVersion", "10.0.4")
.HasAnnotation("Relational:MaxIdentifierLength", 63);
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Avatar")
.HasColumnType("text");
b.Property<DateTime>("CreatedAt")
.HasColumnType("timestamp with time zone");
b.Property<string>("Description")
.HasColumnType("text");
b.Property<bool>("IsImporting")
.HasColumnType("boolean");
b.Property<long>("LastMessageSequenceId")
.HasColumnType("bigint");
b.Property<string>("Name")
.HasColumnType("text");
b.Property<string>("Type")
.IsRequired()
.HasColumnType("text");
b.HasKey("Id");
b.ToTable("Chats", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Folder", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Icon")
.HasColumnType("text");
b.Property<bool>("IsDefault")
.HasColumnType("boolean");
b.Property<string>("Name")
.IsRequired()
.HasColumnType("text");
b.Property<string>("Type")
.IsRequired()
.HasColumnType("text");
b.HasKey("Id");
b.ToTable("Folders", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.UserChatSettings", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<Guid>("ChatId")
.HasColumnType("uuid");
b.Property<string>("FolderIds")
.IsRequired()
.HasColumnType("text");
b.Property<bool>("IsMuted")
.HasColumnType("boolean");
b.Property<Guid>("UserId")
.HasColumnType("uuid");
b.HasKey("Id");
b.HasIndex("UserId", "ChatId")
.IsUnique();
b.ToTable("UserChatSettings", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.OwnsMany("Knot.Contracts.Conversations.Domain.ChatMember", "Members", b1 =>
{
b1.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b1.Property<Guid>("ChatId")
.HasColumnType("uuid");
b1.Property<bool>("IsMuted")
.HasColumnType("boolean");
b1.Property<bool>("IsPinned")
.HasColumnType("boolean");
b1.Property<DateTime>("JoinedAt")
.HasColumnType("timestamp with time zone");
b1.Property<Guid?>("LastDeliveredMessageId")
.HasColumnType("uuid");
b1.Property<Guid?>("LastReadMessageId")
.HasColumnType("uuid");
b1.Property<long>("LastReadSequenceId")
.HasColumnType("bigint");
b1.Property<string>("Role")
.IsRequired()
.HasColumnType("text");
b1.Property<Guid>("UserId")
.HasColumnType("uuid");
b1.HasKey("Id");
b1.HasIndex("ChatId", "UserId")
.IsUnique();
b1.ToTable("ChatMembers", "chats");
b1.WithOwner()
.HasForeignKey("ChatId");
});
b.Navigation("Members");
});
#pragma warning restore 612, 618
}
}
}
@@ -0,0 +1,32 @@
using System;
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace Knot.Modules.Conversations.Migrations
{
/// <inheritdoc />
public partial class AddIsImportingToChat : Migration
{
/// <inheritdoc />
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.AddColumn<bool>(
name: "IsImporting",
schema: "chats",
table: "Chats",
type: "boolean",
nullable: false,
defaultValue: false);
}
/// <inheritdoc />
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropColumn(
name: "IsImporting",
schema: "chats",
table: "Chats");
}
}
}
@@ -0,0 +1,169 @@
// <auto-generated />
using System;
using Knot.Modules.Conversations.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Infrastructure;
using Microsoft.EntityFrameworkCore.Migrations;
using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
#nullable disable
namespace Knot.Modules.Conversations.Migrations
{
[DbContext(typeof(ChatsDbContext))]
[Migration("20260406123444_AddImportJobIdToChat")]
partial class AddImportJobIdToChat
{
/// <inheritdoc />
protected override void BuildTargetModel(ModelBuilder modelBuilder)
{
#pragma warning disable 612, 618
modelBuilder
.HasDefaultSchema("chats")
.HasAnnotation("ProductVersion", "10.0.4")
.HasAnnotation("Relational:MaxIdentifierLength", 63);
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Avatar")
.HasColumnType("text");
b.Property<DateTime>("CreatedAt")
.HasColumnType("timestamp with time zone");
b.Property<string>("Description")
.HasColumnType("text");
b.Property<Guid?>("ImportJobId")
.HasColumnType("uuid");
b.Property<bool>("IsImporting")
.HasColumnType("boolean");
b.Property<long>("LastMessageSequenceId")
.HasColumnType("bigint");
b.Property<string>("Name")
.HasColumnType("text");
b.Property<string>("Type")
.IsRequired()
.HasColumnType("text");
b.HasKey("Id");
b.ToTable("Chats", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Folder", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Icon")
.HasColumnType("text");
b.Property<bool>("IsDefault")
.HasColumnType("boolean");
b.Property<string>("Name")
.IsRequired()
.HasColumnType("text");
b.Property<string>("Type")
.IsRequired()
.HasColumnType("text");
b.HasKey("Id");
b.ToTable("Folders", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.UserChatSettings", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<Guid>("ChatId")
.HasColumnType("uuid");
b.Property<string>("FolderIds")
.IsRequired()
.HasColumnType("text");
b.Property<bool>("IsMuted")
.HasColumnType("boolean");
b.Property<Guid>("UserId")
.HasColumnType("uuid");
b.HasKey("Id");
b.HasIndex("UserId", "ChatId")
.IsUnique();
b.ToTable("UserChatSettings", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.OwnsMany("Knot.Contracts.Conversations.Domain.ChatMember", "Members", b1 =>
{
b1.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b1.Property<Guid>("ChatId")
.HasColumnType("uuid");
b1.Property<bool>("IsMuted")
.HasColumnType("boolean");
b1.Property<bool>("IsPinned")
.HasColumnType("boolean");
b1.Property<DateTime>("JoinedAt")
.HasColumnType("timestamp with time zone");
b1.Property<Guid?>("LastDeliveredMessageId")
.HasColumnType("uuid");
b1.Property<Guid?>("LastReadMessageId")
.HasColumnType("uuid");
b1.Property<long>("LastReadSequenceId")
.HasColumnType("bigint");
b1.Property<string>("Role")
.IsRequired()
.HasColumnType("text");
b1.Property<Guid>("UserId")
.HasColumnType("uuid");
b1.HasKey("Id");
b1.HasIndex("ChatId", "UserId")
.IsUnique();
b1.ToTable("ChatMembers", "chats");
b1.WithOwner()
.HasForeignKey("ChatId");
});
b.Navigation("Members");
});
#pragma warning restore 612, 618
}
}
}
@@ -0,0 +1,31 @@
using System;
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace Knot.Modules.Conversations.Migrations
{
/// <inheritdoc />
public partial class AddImportJobIdToChat : Migration
{
/// <inheritdoc />
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.AddColumn<Guid>(
name: "ImportJobId",
schema: "chats",
table: "Chats",
type: "uuid",
nullable: true);
}
/// <inheritdoc />
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropColumn(
name: "ImportJobId",
schema: "chats",
table: "Chats");
}
}
}
@@ -23,7 +23,7 @@ namespace Knot.Modules.Conversations.Migrations
NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder);
modelBuilder.Entity("Knot.Modules.Conversations.Domain.Chat", b =>
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
@@ -38,6 +38,12 @@ namespace Knot.Modules.Conversations.Migrations
b.Property<string>("Description")
.HasColumnType("text");
b.Property<Guid?>("ImportJobId")
.HasColumnType("uuid");
b.Property<bool>("IsImporting")
.HasColumnType("boolean");
b.Property<long>("LastMessageSequenceId")
.HasColumnType("bigint");
@@ -53,9 +59,61 @@ namespace Knot.Modules.Conversations.Migrations
b.ToTable("Chats", "chats");
});
modelBuilder.Entity("Knot.Modules.Conversations.Domain.Chat", b =>
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Folder", b =>
{
b.OwnsMany("Knot.Modules.Conversations.Domain.ChatMember", "Members", b1 =>
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<string>("Icon")
.HasColumnType("text");
b.Property<bool>("IsDefault")
.HasColumnType("boolean");
b.Property<string>("Name")
.IsRequired()
.HasColumnType("text");
b.Property<string>("Type")
.IsRequired()
.HasColumnType("text");
b.HasKey("Id");
b.ToTable("Folders", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.UserChatSettings", b =>
{
b.Property<Guid>("Id")
.ValueGeneratedOnAdd()
.HasColumnType("uuid");
b.Property<Guid>("ChatId")
.HasColumnType("uuid");
b.Property<string>("FolderIds")
.IsRequired()
.HasColumnType("text");
b.Property<bool>("IsMuted")
.HasColumnType("boolean");
b.Property<Guid>("UserId")
.HasColumnType("uuid");
b.HasKey("Id");
b.HasIndex("UserId", "ChatId")
.IsUnique();
b.ToTable("UserChatSettings", "chats");
});
modelBuilder.Entity("Knot.Contracts.Conversations.Domain.Chat", b =>
{
b.OwnsMany("Knot.Contracts.Conversations.Domain.ChatMember", "Members", b1 =>
{
b1.Property<Guid>("Id")
.ValueGeneratedOnAdd()
@@ -9,7 +9,7 @@ using Knot.Modules.Conversations.Application.Chats.Members;
using Knot.Modules.Conversations.Application.Chats.TogglePin;
using Knot.Modules.Conversations.Application.Chats.Update;
using Knot.Modules.Conversations.Application.DTOs;
using Knot.Modules.Conversations.Domain;
using Knot.Contracts.Conversations.Domain;
using Knot.Shared.Kernel;
using MediatR;
using Microsoft.AspNetCore.Builder;
@@ -90,7 +90,11 @@ public sealed class MessageSentDomainEventHandler : INotificationHandler<Message
storyMediaType = (message as StoryMessage)?.StoryMediaType,
callType = (message as CallMessage)?.CallType,
callStatus = (message as CallMessage)?.CallStatus,
duration = (message as CallMessage)?.Duration
duration = (message as CallMessage)?.Duration,
pollOptions = (message as PollMessage)?.Options.Select(o => new { id = o.Id, text = o.Text, voteCount = o.VoteCount }).ToList(),
pollIsMultipleChoice = (message as PollMessage)?.IsMultipleChoice,
pollIsClosed = (message as PollMessage)?.IsClosed,
pollIsAnonymous = (message as PollMessage)?.IsAnonymous
}, cancellationToken);
}
}
@@ -32,6 +32,18 @@ public sealed class MessageQueryService : IMessageQueryService
public async Task<List<MessageInfo>> GetOrphanedMessagesAsync(HashSet<Guid> activeChatIds, CancellationToken cancellationToken)
{
if (activeChatIds == null || activeChatIds.Count == 0)
{
var allMessages = await _messages.Find(_ => true).ToListAsync(cancellationToken);
return allMessages.Select(m => new MessageInfo(
m.Id,
m.ChatId,
m.State.HasFlag(MessageState.IsDeleted),
m is MediaMessage mm && mm.Media.Any() ? mm.Media.First().Url : null,
m is MediaMessage mm2 && mm2.Media.Any() ? mm2.Media.Select(media => new Contracts.Messaging.Application.Abstractions.MediaInfo(media.Url)).ToList() : null
)).ToList();
}
var builder = Builders<Message>.Filter;
var inFilter = builder.In(m => m.ChatId, activeChatIds);
var filter = builder.Not(inFilter);
@@ -62,22 +62,69 @@ public sealed class MessageRepository : IMessageRepository
.FirstOrDefaultAsync(cancellationToken);
}
public async Task<List<Message>> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, int limit, CancellationToken cancellationToken)
public async Task<List<Message>> GetPinnedMessagesAsync(Guid chatId, CancellationToken cancellationToken)
{
var builder = Builders<Message>.Filter;
var filter = builder.And(
builder.Eq(m => m.ChatId, chatId),
builder.BitsAnySet(m => m.State, (long)MessageState.IsPinned)
);
return await _messages.Find(filter)
.SortByDescending(m => m.CreatedAt)
.ToListAsync(cancellationToken);
}
public async Task<List<Message>> GetChatMessagesCursorAsync(Guid chatId, DateTime? cursor, long? sequenceId, int limit, CancellationToken cancellationToken)
{
var builder = Builders<Message>.Filter;
var filter = builder.Eq(m => m.ChatId, chatId);
if (cursor.HasValue)
if (sequenceId.HasValue)
{
filter &= builder.Lt(m => m.SequenceId, sequenceId.Value);
}
else if (cursor.HasValue)
{
filter &= builder.Lt(m => m.CreatedAt, cursor.Value);
}
return await _messages.Find(filter)
.SortByDescending(m => m.CreatedAt)
.SortByDescending(m => m.SequenceId)
.Limit(limit)
.ToListAsync(cancellationToken);
}
public async Task<List<Message>> GetChatMessagesAroundAsync(Guid chatId, long sequenceId, int limit, CancellationToken cancellationToken)
{
var builder = Builders<Message>.Filter;
// Target message
var targetFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Eq(m => m.SequenceId, sequenceId));
var targetMsg = await _messages.Find(targetFilter).FirstOrDefaultAsync(cancellationToken);
// Older messages
var olderFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Lt(m => m.SequenceId, sequenceId));
var older = await _messages.Find(olderFilter)
.SortByDescending(m => m.SequenceId)
.Limit(limit / 2)
.ToListAsync(cancellationToken);
// Newer messages
var newerFilter = builder.And(builder.Eq(m => m.ChatId, chatId), builder.Gt(m => m.SequenceId, sequenceId));
var newer = await _messages.Find(newerFilter)
.SortBy(m => m.SequenceId)
.Limit(limit / 2)
.ToListAsync(cancellationToken);
var result = new List<Message>();
result.AddRange(older);
if (targetMsg != null) result.Add(targetMsg);
result.AddRange(newer);
return result.OrderBy(m => m.SequenceId).ToList();
}
public async Task<List<Message>> SearchMessagesAsync(string query, Guid? chatId, Guid requestingUserId, CancellationToken cancellationToken)
{
// Not ideal for SQL/Mongo combination but keeping the signature
@@ -2,5 +2,9 @@ using System;
namespace Knot.Modules.Profiles.Application.Profiles.DTOs;
public record UpdateProfileRequest(string? DisplayName, string? Bio, DateTime? Birthday);
public record UpdateProfileRequest(
string? DisplayName,
string? Bio,
DateTime? Birthday,
bool? IsInvisible);
public record UpdateSettingsRequest(bool? HideStoryViews);
@@ -11,7 +11,8 @@ public sealed record UpdateProfileCommand(
Guid UserId,
string? DisplayName,
string? Bio,
DateTime? Birthday) : ICommand<UserProfileDto>;
DateTime? Birthday,
bool? IsInvisible) : ICommand<UserProfileDto>;
internal sealed class UpdateProfileCommandHandler : ICommandHandler<UpdateProfileCommand, UserProfileDto>
{
@@ -28,9 +29,13 @@ internal sealed class UpdateProfileCommandHandler : ICommandHandler<UpdateProfil
if (profile is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
if ((request.Bio?.Length ?? 0) > 200)
return Result.Failure<UserProfileDto>(ProfilesErrors.BioTooLong);
profile.DisplayName = request.DisplayName ?? profile.DisplayName;
profile.About = request.Bio;
profile.Birthday = request.Birthday;
profile.IsInvisible = request.IsInvisible ?? profile.IsInvisible;
var result = await _repository.UpdateAsync(profile, cancellationToken);
if (result.IsFailure)
@@ -0,0 +1,21 @@
using Knot.Contracts.Profiles.Application.DTOs;
using Knot.Contracts.Profiles.Domain;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Profiles.Application.Statuses;
public sealed record ClearUserStatusCommand(Guid UserId) : IRequest<Result<UserProfileDto>>;
public sealed class ClearUserStatusCommandHandler : IRequestHandler<ClearUserStatusCommand, Result<UserProfileDto>>
{
private readonly IProfileStatusWriter _writer;
public ClearUserStatusCommandHandler(IProfileStatusWriter writer)
{
_writer = writer;
}
public Task<Result<UserProfileDto>> Handle(ClearUserStatusCommand request, CancellationToken cancellationToken)
=> _writer.ClearAsync(request.UserId, cancellationToken);
}
@@ -0,0 +1,12 @@
using Knot.Contracts.Profiles.Application.DTOs;
using MediatR;
namespace Knot.Modules.Profiles.Application.Statuses;
public sealed record GetStatusPresetsQuery : IRequest<IReadOnlyList<StatusPresetDto>>;
public sealed class GetStatusPresetsQueryHandler : IRequestHandler<GetStatusPresetsQuery, IReadOnlyList<StatusPresetDto>>
{
public Task<IReadOnlyList<StatusPresetDto>> Handle(GetStatusPresetsQuery request, CancellationToken cancellationToken)
=> Task.FromResult(StatusPresetCatalog.All);
}
@@ -0,0 +1,53 @@
using Knot.Contracts.Profiles.Application.DTOs;
using Knot.Contracts.Profiles.Domain;
using Knot.Shared.Kernel;
using MediatR;
namespace Knot.Modules.Profiles.Application.Statuses;
public sealed record SetCustomUserStatusCommand(
Guid UserId,
string? Emoji,
string? Text,
DateTime? ExpiresAt,
string? PresetKey) : IRequest<Result<UserProfileDto>>;
public sealed class SetCustomUserStatusCommandHandler : IRequestHandler<SetCustomUserStatusCommand, Result<UserProfileDto>>
{
private readonly IProfileStatusWriter _writer;
public SetCustomUserStatusCommandHandler(IProfileStatusWriter writer)
{
_writer = writer;
}
public async Task<Result<UserProfileDto>> Handle(SetCustomUserStatusCommand request, CancellationToken cancellationToken)
{
var key = request.PresetKey?.Trim();
if (!string.IsNullOrEmpty(key) && key.Equals("online", StringComparison.OrdinalIgnoreCase))
return await _writer.ClearAsync(request.UserId, cancellationToken);
var emoji = request.Emoji?.Trim() ?? string.Empty;
var text = request.Text?.Trim() ?? string.Empty;
if (!string.IsNullOrEmpty(key))
{
var preset = StatusPresetCatalog.Find(key);
if (preset is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.InvalidPreset);
if (string.IsNullOrEmpty(emoji))
emoji = preset.Emoji;
if (string.IsNullOrEmpty(text))
text = preset.TextRu;
}
if (string.IsNullOrEmpty(emoji) && string.IsNullOrEmpty(text))
return Result.Failure<UserProfileDto>(ProfilesErrors.StatusEmpty);
if (text.Length > 50)
return Result.Failure<UserProfileDto>(ProfilesErrors.StatusTextTooLong);
var presetForWriter = string.IsNullOrWhiteSpace(key) ? null : key;
return await _writer.SetCustomAsync(request.UserId, emoji, text, request.ExpiresAt, presetForWriter, cancellationToken);
}
}
@@ -0,0 +1,24 @@
using Knot.Contracts.Profiles.Application.DTOs;
namespace Knot.Modules.Profiles.Application.Statuses;
public static class StatusPresetCatalog
{
private static readonly StatusPresetDto[] Items =
[
new() { Id = "online", Emoji = "", TextRu = "Онлайн (стандарт)", TextEn = "Online (standard)" },
new() { Id = "away", Emoji = "🕒", TextRu = "В отъезде", TextEn = "Away" },
new() { Id = "dnd", Emoji = "⛔", TextRu = "Не беспокоить", TextEn = "Do not disturb" },
new() { Id = "sick", Emoji = "🤒", TextRu = "Болен", TextEn = "Sick" },
new() { Id = "angry", Emoji = "💢", TextRu = "Злой", TextEn = "Angry" }
];
public static IReadOnlyList<StatusPresetDto> All => Items;
public static StatusPresetDto? Find(string presetKey)
{
if (string.IsNullOrWhiteSpace(presetKey))
return null;
return Items.FirstOrDefault(x => x.Id.Equals(presetKey.Trim(), StringComparison.OrdinalIgnoreCase));
}
}
@@ -11,10 +11,11 @@ public static class DependencyInjection
public static IServiceCollection AddProfilesModule(this IServiceCollection services, IConfiguration configuration)
{
services.AddScoped<Knot.Contracts.Profiles.Domain.IProfileRepository, ProfileRepository>();
services.AddScoped<Knot.Contracts.Profiles.Domain.IUserStatusRepository, UserStatusMongoRepository>();
services.AddScoped<Knot.Contracts.Profiles.Domain.IProfileStatusWriter, ProfileStatusWriter>();
services.AddScoped<Knot.Contracts.Profiles.Domain.IProfilesUnitOfWork, ProfilesUnitOfWork>();
services.AddScoped<Knot.Contracts.Profiles.Domain.IAvatarStorageService, AvatarStorageService>();
// MongoDB Registration
var mongoConnection = configuration.GetConnectionString("MongoConnection")
?? configuration["MONGO_URL"]
?? "mongodb://mongo:27017";
@@ -3,10 +3,6 @@ using MongoDB.Bson.Serialization.Attributes;
namespace Knot.Modules.Profiles.Domain;
/// <summary>
/// MongoDB-документ профиля пользователя.
/// Id совпадает с UserId из модуля Auth (Postgres).
/// </summary>
public sealed class ProfileDocument
{
[BsonId]
@@ -19,12 +15,23 @@ public sealed class ProfileDocument
public string? Bio { get; private set; }
public string? StatusText { get; private set; }
public string? StatusEmoji { get; private set; }
public DateTime? StatusExpiresAt { get; private set; }
[BsonRepresentation(BsonType.String)]
public Guid? CurrentStatusId { get; private set; }
public string? AvatarUrl { get; private set; }
public DateTime? Birthday { get; private set; }
public bool HideStoryViews { get; private set; }
public bool IsInvisible { get; private set; }
public bool IsBanned { get; private set; }
public bool IsDeleted { get; private set; }
@@ -65,6 +72,26 @@ public sealed class ProfileDocument
public void UpdateSettings(bool hideStoryViews)
=> HideStoryViews = hideStoryViews;
public void UpdateInvisible(bool isInvisible)
=> IsInvisible = isInvisible;
public void UpdateStatus(string? statusText, string? statusEmoji, DateTime? statusExpiresAt)
{
StatusText = statusText;
StatusEmoji = statusEmoji;
StatusExpiresAt = statusExpiresAt;
}
public void SetCurrentStatusId(Guid? statusId)
=> CurrentStatusId = statusId;
public void ClearMoodStatusFields()
{
StatusText = null;
StatusEmoji = null;
StatusExpiresAt = null;
}
public void UpdateStatus(bool isBanned, bool isDeleted)
{
IsBanned = isBanned;
@@ -0,0 +1,26 @@
using MongoDB.Bson;
using MongoDB.Bson.Serialization.Attributes;
namespace Knot.Modules.Profiles.Domain;
public sealed class UserStatusDocument
{
[BsonId]
[BsonRepresentation(BsonType.String)]
public Guid Id { get; set; }
[BsonRepresentation(BsonType.String)]
public Guid UserId { get; set; }
public string Type { get; set; } = "Custom";
public string Emoji { get; set; } = "";
public string Text { get; set; } = "";
public DateTime CreatedAt { get; set; }
public DateTime? ExpiresAt { get; set; }
public BsonDocument Metadata { get; set; } = new();
}
@@ -10,29 +10,39 @@ namespace Knot.Modules.Profiles.Infrastructure.Database;
internal class ProfileRepository : IProfileRepository
{
private readonly IMongoCollection<ProfileDocument> _profiles;
private readonly IUserStatusRepository _statuses;
public ProfileRepository(IMongoDatabase database)
public ProfileRepository(IMongoDatabase database, IUserStatusRepository statuses)
{
_profiles = database.GetCollection<ProfileDocument>("profiles");
_statuses = statuses;
}
public async Task<UserProfileDto?> GetAsync(Guid userId, CancellationToken ct = default)
{
var profile = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(ct);
return profile?.ToDto();
return profile is null ? null : await profile.ToDtoAsync(_statuses, ct);
}
public async Task<UserProfileDto?> GetByUsernameAsync(string username, CancellationToken ct = default)
{
var profile = await _profiles.Find(p => p.Username == username).FirstOrDefaultAsync(ct);
return profile?.ToDto();
return profile is null ? null : await profile.ToDtoAsync(_statuses, ct);
}
public async Task<List<UserProfileDto>> GetAsync(IEnumerable<Guid> userIds, CancellationToken ct = default)
{
var ids = userIds.ToList();
var profiles = await _profiles.Find(p => ids.Contains(p.Id)).ToListAsync(ct);
return profiles.Select(p => p.ToDto()).ToList();
var statusIds = profiles
.Where(p => p.CurrentStatusId.HasValue)
.Select(p => p.CurrentStatusId!.Value)
.Distinct()
.ToList();
var map = statusIds.Count > 0
? await _statuses.GetByIdsAsync(statusIds, ct)
: new Dictionary<Guid, UserStatusDto>();
return profiles.Select(p => ProfileMappings.ToDtoWithStatusMap(p, map)).ToList();
}
public async Task<List<UserProfileDto>> SearchAsync(string query, int limit = 20, CancellationToken ct = default)
@@ -57,7 +67,16 @@ internal class ProfileRepository : IProfileRepository
var combinedFilter = Builders<ProfileDocument>.Filter.And(baseFilter, searchFilter);
docs = await _profiles.Find(combinedFilter).Limit(limit).ToListAsync(ct);
}
return docs.Select(p => p.ToDto()).ToList();
var statusIds = docs
.Where(p => p.CurrentStatusId.HasValue)
.Select(p => p.CurrentStatusId!.Value)
.Distinct()
.ToList();
var map = statusIds.Count > 0
? await _statuses.GetByIdsAsync(statusIds, ct)
: new Dictionary<Guid, UserStatusDto>();
return docs.Select(p => ProfileMappings.ToDtoWithStatusMap(p, map)).ToList();
}
public async Task<Result<UserProfileDto>> CreateAsync(UserProfileDto dto, CancellationToken ct = default)
@@ -65,7 +84,7 @@ internal class ProfileRepository : IProfileRepository
var profile = ProfileDocument.Create(dto.UserId, dto.Username ?? string.Empty, dto.DisplayName ?? string.Empty, dto.About);
await _profiles.InsertOneAsync(profile, null, ct);
return Result.Success(profile.ToDto());
return Result.Success(await profile.ToDtoAsync(_statuses, ct));
}
public async Task<Result<UserProfileDto>> UpdateAsync(UserProfileDto dto, CancellationToken ct = default)
@@ -80,9 +99,13 @@ internal class ProfileRepository : IProfileRepository
dto.Birthday);
document.UpdateAvatar(dto.Avatar);
document.UpdateInvisible(dto.IsInvisible);
await _profiles.ReplaceOneAsync(p => p.Id == dto.UserId, document, cancellationToken: ct);
return Result.Success(document.ToDto());
var fresh = await _profiles.Find(p => p.Id == dto.UserId).FirstOrDefaultAsync(ct);
if (fresh is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
return Result.Success(await fresh.ToDtoAsync(_statuses, ct));
}
public async Task<Result> UpdateStatusAsync(Guid userId, bool isBanned, bool isDeleted, CancellationToken ct = default)
@@ -0,0 +1,73 @@
using Knot.Contracts.Profiles.Application.DTOs;
using Knot.Contracts.Profiles.Domain;
using Knot.Modules.Profiles.Domain;
using Knot.Modules.Profiles.Infrastructure.Mappings;
using Knot.Shared.Kernel;
using MongoDB.Driver;
namespace Knot.Modules.Profiles.Infrastructure.Database;
internal sealed class ProfileStatusWriter : IProfileStatusWriter
{
private readonly IMongoCollection<ProfileDocument> _profiles;
private readonly IUserStatusRepository _statuses;
public ProfileStatusWriter(IMongoDatabase database, IUserStatusRepository statuses)
{
_profiles = database.GetCollection<ProfileDocument>("profiles");
_statuses = statuses;
}
public async Task<Result<UserProfileDto>> SetCustomAsync(Guid userId, string emoji, string text, DateTime? expiresAt, string? presetKey, CancellationToken cancellationToken = default)
{
var profile = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(cancellationToken);
if (profile is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
if (profile.CurrentStatusId is Guid oldId)
await _statuses.DeleteAsync(oldId, cancellationToken);
var newId = Guid.NewGuid();
var type = string.IsNullOrWhiteSpace(presetKey) ? "Custom" : "Preset";
var dto = new UserStatusDto
{
Id = newId,
Type = type,
Emoji = emoji,
Text = text,
CreatedAt = DateTime.UtcNow,
ExpiresAt = expiresAt
};
await _statuses.InsertAsync(dto, userId, cancellationToken);
profile.SetCurrentStatusId(newId);
profile.ClearMoodStatusFields();
await _profiles.ReplaceOneAsync(p => p.Id == userId, profile, cancellationToken: cancellationToken);
var fresh = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(cancellationToken);
if (fresh is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
return Result.Success(await fresh.ToDtoAsync(_statuses, cancellationToken));
}
public async Task<Result<UserProfileDto>> ClearAsync(Guid userId, CancellationToken cancellationToken = default)
{
var profile = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(cancellationToken);
if (profile is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
if (profile.CurrentStatusId is Guid oldId)
await _statuses.DeleteAsync(oldId, cancellationToken);
profile.SetCurrentStatusId(null);
profile.ClearMoodStatusFields();
await _profiles.ReplaceOneAsync(p => p.Id == userId, profile, cancellationToken: cancellationToken);
var fresh = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(cancellationToken);
if (fresh is null)
return Result.Failure<UserProfileDto>(ProfilesErrors.ProfileNotFound);
return Result.Success(await fresh.ToDtoAsync(_statuses, cancellationToken));
}
}
@@ -0,0 +1,62 @@
using Knot.Contracts.Profiles.Application.DTOs;
using Knot.Contracts.Profiles.Domain;
using Knot.Modules.Profiles.Domain;
using MongoDB.Bson;
using MongoDB.Driver;
namespace Knot.Modules.Profiles.Infrastructure.Database;
internal sealed class UserStatusMongoRepository : IUserStatusRepository
{
private readonly IMongoCollection<UserStatusDocument> _collection;
public UserStatusMongoRepository(IMongoDatabase database)
{
_collection = database.GetCollection<UserStatusDocument>("user_statuses");
}
public async Task<UserStatusDto?> GetByIdAsync(Guid id, CancellationToken cancellationToken = default)
{
var doc = await _collection.Find(x => x.Id == id).FirstOrDefaultAsync(cancellationToken);
return doc is null ? null : ToDto(doc);
}
public async Task<IReadOnlyDictionary<Guid, UserStatusDto>> GetByIdsAsync(IEnumerable<Guid> ids, CancellationToken cancellationToken = default)
{
var idList = ids.Distinct().ToList();
if (idList.Count == 0)
return new Dictionary<Guid, UserStatusDto>();
var docs = await _collection.Find(x => idList.Contains(x.Id)).ToListAsync(cancellationToken);
return docs.ToDictionary(d => d.Id, ToDto);
}
public async Task InsertAsync(UserStatusDto dto, Guid userId, CancellationToken cancellationToken = default)
{
var doc = new UserStatusDocument
{
Id = dto.Id,
UserId = userId,
Type = dto.Type,
Emoji = dto.Emoji,
Text = dto.Text,
CreatedAt = dto.CreatedAt,
ExpiresAt = dto.ExpiresAt,
Metadata = new BsonDocument()
};
await _collection.InsertOneAsync(doc, cancellationToken: cancellationToken);
}
public Task DeleteAsync(Guid id, CancellationToken cancellationToken = default)
=> _collection.DeleteOneAsync(x => x.Id == id, cancellationToken);
private static UserStatusDto ToDto(UserStatusDocument d) => new()
{
Id = d.Id,
Type = d.Type,
Emoji = d.Emoji,
Text = d.Text,
CreatedAt = d.CreatedAt,
ExpiresAt = d.ExpiresAt
};
}

Some files were not shown because too many files have changed in this diff Show More