Рабочий чат

This commit is contained in:
Халимов Рустам
2026-04-01 23:18:55 +03:00
parent 249c344df8
commit 4ae7dd60ce
42 changed files with 1016 additions and 184 deletions

View File

@@ -0,0 +1,31 @@
using Knot.Shared.Kernel;
using Knot.Contracts.Auth.Application.Abstractions;
using Knot.Contracts.Auth.Domain;
using Microsoft.EntityFrameworkCore;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
namespace Knot.Modules.Auth.Application.Users;
internal sealed class GetUsersExistenceQueryHandler : IQueryHandler<GetUsersExistenceQuery, List<Guid>>
{
private readonly IAuthDbContext _context;
public GetUsersExistenceQueryHandler(IAuthDbContext context)
{
_context = context;
}
public async Task<Result<List<Guid>>> Handle(GetUsersExistenceQuery request, CancellationToken cancellationToken)
{
var existingIds = await _context.Set<Knot.Modules.Auth.Domain.User>()
.Where(u => request.UserIds.Contains(u.Id))
.Select(u => u.Id)
.ToListAsync(cancellationToken);
return Result.Success(existingIds);
}
}

View File

@@ -28,8 +28,14 @@ public sealed class User : AggregateRoot<Guid>
public string? RefreshToken { get; private set; }
public DateTime? BannedUntil { get; private set; }
public void Ban() => IsBanned = true;
public void Unban() => IsBanned = false;
public void Ban() {
IsBanned = true;
RaiseDomainEvent(new UserBannedDomainEvent(Id, true));
}
public void Unban() {
IsBanned = false;
RaiseDomainEvent(new UserBannedDomainEvent(Id, false));
}
public void SetOnline(bool isOnline, DateTime? lastSeen = null)
{

View File

@@ -225,18 +225,31 @@ public sealed class ChatHub : Hub
[HubMethodName("friend_request")]
public async Task FriendRequest(FriendSignalRequest request)
{
if (request == null || string.IsNullOrEmpty(request.FriendId))
{
_logger.LogWarning("FriendRequest called with null request or empty FriendId");
return;
}
_logger.LogInformation("Signaling friend_request_received to {FriendId} from {UserId}", request.FriendId, _userContext.UserId);
await SendToUserAsync(request.FriendId, "friend_request_received", new { userId = _userContext.UserId });
}
[HubMethodName("friend_accepted")]
public async Task FriendAccepted(FriendSignalRequest request)
{
if (request == null || string.IsNullOrEmpty(request.FriendId)) return;
_logger.LogInformation("Signaling friend_request_accepted to {FriendId} from {UserId}", request.FriendId, _userContext.UserId);
await SendToUserAsync(request.FriendId, "friend_request_accepted", new { userId = _userContext.UserId });
}
[HubMethodName("friend_removed")]
public async Task FriendRemoved(FriendSignalRequest request)
{
if (request == null || string.IsNullOrEmpty(request.FriendId)) return;
_logger.LogInformation("Signaling friend_removed_notify to {FriendId} from {UserId}", request.FriendId, _userContext.UserId);
await SendToUserAsync(request.FriendId, "friend_removed_notify", new { userId = _userContext.UserId });
}
@@ -246,15 +259,26 @@ public sealed class ChatHub : Hub
private async Task SendToUserAsync(string targetUserId, string method, object payload)
{
if (string.IsNullOrEmpty(targetUserId))
{
_logger.LogWarning("SendToUserAsync called with null or empty targetUserId");
return;
}
if (_userConnections.TryGetValue(targetUserId, out var connectionIds))
{
string[] ids;
lock (connectionIds) { ids = connectionIds.ToArray(); }
_logger.LogDebug("Sending {Method} to user {TargetUserId} ({ConnectionCount} connections)", method, targetUserId, ids.Length);
foreach (var connId in ids)
{
await Clients.Client(connId).SendAsync(method, payload);
}
}
else
{
_logger.LogDebug("User {TargetUserId} not online, skipping {Method} signal", targetUserId, method);
}
}
[HubMethodName("call_offer")]

View File

@@ -0,0 +1,31 @@
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Events;
using Knot.Contracts.Profiles.Domain;
using Knot.Modules.Profiles.Domain;
using MediatR;
using System.Threading;
using System.Threading.Tasks;
namespace Knot.Modules.Profiles.Application.Profiles.Integration;
internal sealed class UserStatusChangedHandler :
INotificationHandler<UserBannedDomainEvent>,
INotificationHandler<UserDeletedDomainEvent>
{
private readonly IProfileRepository _repository;
public UserStatusChangedHandler(IProfileRepository repository)
{
_repository = repository;
}
public async Task Handle(UserBannedDomainEvent notification, CancellationToken cancellationToken)
{
await _repository.UpdateStatusAsync(notification.UserId, notification.IsBanned, false, cancellationToken);
}
public async Task Handle(UserDeletedDomainEvent notification, CancellationToken cancellationToken)
{
await _repository.DeleteAsync(notification.UserId, cancellationToken);
}
}

View File

@@ -5,23 +5,69 @@ using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using MediatR;
using System;
using Knot.Contracts.Relations.Application.Contacts;
using Knot.Contracts.Auth.Application.Abstractions;
namespace Knot.Modules.Profiles.Application.Profiles.Search;
public sealed record SearchProfilesQuery(string Query) : IQuery<List<UserProfileDto>>;
public sealed record SearchProfilesQuery(string Query, Guid UserId) : IQuery<List<UserProfileDto>>;
internal sealed class SearchProfilesQueryHandler : IQueryHandler<SearchProfilesQuery, List<UserProfileDto>>
{
private readonly IProfileRepository _repository;
private readonly ISender _sender;
public SearchProfilesQueryHandler(IProfileRepository repository)
public SearchProfilesQueryHandler(IProfileRepository repository, ISender sender)
{
_repository = repository;
_sender = sender;
}
public async Task<Result<List<UserProfileDto>>> Handle(SearchProfilesQuery request, CancellationToken cancellationToken)
{
// 1. Fetch candidates from MongoDB (which might have orphans)
var profiles = await _repository.SearchAsync(request.Query, 20, cancellationToken);
return Result.Success(profiles);
var candidates = profiles.Where(p => p.UserId != request.UserId).ToList();
if (!candidates.Any()) return Result.Success(new List<UserProfileDto>());
var candidateIds = candidates.Select(p => p.UserId).ToList();
// 2. Validate existence in Auth module (to filter out orphans from deleted accounts)
var existingIds = new List<Guid>();
try {
var existenceResult = await _sender.Send(new GetUsersExistenceQuery(candidateIds), cancellationToken);
if (existenceResult.IsSuccess) existingIds = existenceResult.Value;
} catch {
// If Auth module is not available, we assume all exist to avoid empty results
existingIds = candidateIds;
}
// 3. Remove orphans (and physically delete them from Mongo if they don't exist in Auth)
var validCandidates = candidates.Where(p => existingIds.Contains(p.UserId)).ToList();
var orphanIds = candidateIds.Except(existingIds).ToList();
foreach (var orphanId in orphanIds)
{
// Background cleanup (fire and forget or just do it since it's only a few)
_ = _repository.DeleteAsync(orphanId, CancellationToken.None);
}
if (!validCandidates.Any()) return Result.Success(new List<UserProfileDto>());
// 4. Check for blocked users among valid candidates
var blockedIds = new List<Guid>();
try {
var blockedResult = await _sender.Send(new CheckBlockedStatusQuery(request.UserId, validCandidates.Select(v => v.UserId).ToList()), cancellationToken);
if (blockedResult.IsSuccess) blockedIds = blockedResult.Value;
} catch { }
// final filter
var filtered = validCandidates
.Where(p => !blockedIds.Contains(p.UserId))
.ToList();
return Result.Success(filtered);
}
}

View File

@@ -25,6 +25,10 @@ public sealed class ProfileDocument
public bool HideStoryViews { get; private set; }
public bool IsBanned { get; private set; }
public bool IsDeleted { get; private set; }
public DateTime CreatedAt { get; private set; }
#pragma warning disable CS8618
@@ -37,6 +41,8 @@ public sealed class ProfileDocument
Username = username;
DisplayName = displayName;
Bio = bio;
IsBanned = false;
IsDeleted = false;
CreatedAt = DateTime.UtcNow;
}
@@ -58,4 +64,10 @@ public sealed class ProfileDocument
public void UpdateSettings(bool hideStoryViews)
=> HideStoryViews = hideStoryViews;
public void UpdateStatus(bool isBanned, bool isDeleted)
{
IsBanned = isBanned;
IsDeleted = isDeleted;
}
}

View File

@@ -43,18 +43,25 @@ internal class ProfileRepository : IProfileRepository
public async Task<List<UserProfileDto>> SearchAsync(string query, int limit = 20, CancellationToken ct = default)
{
var baseFilter = Builders<ProfileDocument>.Filter.And(
Builders<ProfileDocument>.Filter.Ne(p => p.IsBanned, true),
Builders<ProfileDocument>.Filter.Ne(p => p.IsDeleted, true)
);
List<ProfileDocument> docs;
if (string.IsNullOrWhiteSpace(query))
{
docs = await _profiles.Find(_ => true).Limit(limit).ToListAsync(ct);
docs = await _profiles.Find(baseFilter).Limit(limit).ToListAsync(ct);
}
else
{
var filter = Builders<ProfileDocument>.Filter.Or(
var searchFilter = Builders<ProfileDocument>.Filter.Or(
Builders<ProfileDocument>.Filter.Regex(p => p.Username, new BsonRegularExpression(query, "i")),
Builders<ProfileDocument>.Filter.Regex(p => p.DisplayName, new BsonRegularExpression(query, "i"))
);
docs = await _profiles.Find(filter).Limit(limit).ToListAsync(ct);
var combinedFilter = Builders<ProfileDocument>.Filter.And(baseFilter, searchFilter);
docs = await _profiles.Find(combinedFilter).Limit(limit).ToListAsync(ct);
}
return docs.Select(p => p.ToDto()).ToList();
}
@@ -82,6 +89,16 @@ internal class ProfileRepository : IProfileRepository
return Result.Success(profile.ToDto());
}
public async Task<Result> UpdateStatusAsync(Guid userId, bool isBanned, bool isDeleted, CancellationToken ct = default)
{
var profile = await _profiles.Find(p => p.Id == userId).FirstOrDefaultAsync(ct);
if (profile is null) return Result.Failure(ProfilesErrors.ProfileNotFound);
profile.UpdateStatus(isBanned, isDeleted);
await _profiles.ReplaceOneAsync(p => p.Id == userId, profile, cancellationToken: ct);
return Result.Success();
}
public async Task<Result> DeleteAsync(Guid userId, CancellationToken ct = default)
{
var result = await _profiles.DeleteOneAsync(p => p.Id == userId, ct);

View File

@@ -10,6 +10,8 @@
<ProjectReference Include="..\..\Shared\Knot.Shared.Kernel\Knot.Shared.Kernel.csproj" />
<ProjectReference Include="..\..\Shared\Knot.Shared.Infrastructure\Knot.Shared.Infrastructure.csproj" />
<ProjectReference Include="..\..\Contracts\Profiles\Knot.Contracts.Profiles.csproj" />
<ProjectReference Include="..\..\Contracts\Relations\Knot.Contracts.Relations.csproj" />
<ProjectReference Include="..\..\Contracts\Auth\Knot.Contracts.Auth.csproj" />
</ItemGroup>
<ItemGroup>

View File

@@ -18,9 +18,9 @@ public static class ProfilesEndpoints
{
var group = app.MapGroup("api/profiles").RequireAuthorization();
group.MapGet("search", async ([FromQuery] string q, ISender sender, CancellationToken ct) =>
group.MapGet("search", async ([FromQuery] string q, ISender sender, IUserContext userContext, CancellationToken ct) =>
{
var result = await sender.Send(new SearchProfilesQuery(q), ct);
var result = await sender.Send(new SearchProfilesQuery(q, userContext.UserId), ct);
return result.IsSuccess ? Results.Ok(result.Value) : Results.BadRequest(result.Error);
});

View File

@@ -0,0 +1,36 @@
using Knot.Shared.Kernel;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using Knot.Modules.Relations.Application.Abstractions;
using Knot.Modules.Relations.Domain;
using Knot.Contracts.Relations.Application.Contacts;
namespace Knot.Modules.Relations.Application.Contacts;
internal sealed class CheckBlockedStatusQueryHandler : IQueryHandler<CheckBlockedStatusQuery, List<Guid>>
{
private readonly IContactsDbContext _context;
public CheckBlockedStatusQueryHandler(IContactsDbContext context)
{
_context = context;
}
public async Task<Result<List<Guid>>> Handle(CheckBlockedStatusQuery request, CancellationToken cancellationToken)
{
// Find which candidate IDs are in a "Blocked" relationship with the current user.
// We check both directions (currentUser blocks candidate OR candidate blocks currentUser).
var blockedIds = await _context.Contacts
.Where(c => (c.UserId == request.UserId || c.ContactId == request.UserId)
&& c.Status == ContactStatus.Blocked
&& (request.CandidateIds.Contains(c.UserId) || request.CandidateIds.Contains(c.ContactId)))
.Select(c => c.UserId == request.UserId ? c.ContactId : c.UserId)
.ToListAsync(cancellationToken);
return Result.Success(blockedIds);
}
}

View File

@@ -22,9 +22,11 @@ internal sealed class DeclineContactRequestCommandHandler : ICommandHandler<Decl
public async Task<Result<bool>> Handle(DeclineContactRequestCommand request, CancellationToken cancellationToken)
{
var contact = await _context.Contacts.FirstOrDefaultAsync(c => c.Id == request.RequestId, cancellationToken);
if (contact == null || contact.ContactId != request.UserId)
// Allow BOTH receiver (to decline) AND sender (to cancel)
if (contact == null || (contact.ContactId != request.UserId && contact.UserId != request.UserId))
{
return Result.Failure<bool>(new Error("Contacts.NotFound", "Contact request not found."));
return Result.Failure<bool>(new Error("Contacts.NotFound", "Contact request not found or you don't have permission."));
}
contact.Decline();

View File

@@ -0,0 +1,33 @@
using Knot.Shared.Kernel;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using Knot.Modules.Relations.Application.Abstractions;
using Knot.Modules.Relations.Domain;
using Knot.Contracts.Relations.Application.Contacts;
namespace Knot.Modules.Relations.Application.Contacts;
internal sealed class GetBlockedUserIdsQueryHandler : IQueryHandler<GetBlockedUserIdsQuery, List<Guid>>
{
private readonly IContactsDbContext _context;
public GetBlockedUserIdsQueryHandler(IContactsDbContext context)
{
_context = context;
}
public async Task<Result<List<Guid>>> Handle(GetBlockedUserIdsQuery request, CancellationToken cancellationToken)
{
var blockedIds = await _context.Contacts
.Where(c => (c.UserId == request.UserId || c.ContactId == request.UserId) && c.Status == ContactStatus.Blocked)
.Select(c => c.UserId == request.UserId ? c.ContactId : c.UserId)
.ToListAsync(cancellationToken);
return Result.Success(blockedIds);
}
}

View File

@@ -0,0 +1,75 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Relations.Domain;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
using Knot.Modules.Relations.Application.Abstractions;
namespace Knot.Modules.Relations.Application.Contacts;
public record ContactUserDto(Guid Id, string Username, string DisplayName, string Avatar);
public record ContactRequestDto(Guid Id, ContactUserDto User, DateTime CreatedAt, bool IsOutgoing);
public record GetContactRequestsQuery(Guid UserId) : IQuery<List<ContactRequestDto>>;
internal sealed class GetContactRequestsQueryHandler : IQueryHandler<GetContactRequestsQuery, List<ContactRequestDto>>
{
private readonly IContactsDbContext _context;
public GetContactRequestsQueryHandler(IContactsDbContext context)
{
_context = context;
}
public async Task<Result<List<ContactRequestDto>>> Handle(GetContactRequestsQuery request, CancellationToken cancellationToken)
{
// 1. Fetch ALL pending requests where current user is either sender or receiver
var contacts = await _context.Contacts
.Where(c => (c.ContactId == request.UserId || c.UserId == request.UserId) && c.Status == ContactStatus.Pending)
.ToListAsync(cancellationToken);
// 2. Fetch all unique IDs for users we need replicas for
var userIds = contacts
.Select(c => c.UserId == request.UserId ? c.ContactId : c.UserId)
.Distinct()
.ToList();
// 3. Fetch replicas
var replicas = await _context.UserReplicas
.Where(r => userIds.Contains(r.Id))
.ToDictionaryAsync(r => r.Id, cancellationToken);
// 4. Transform into DTOs
var result = contacts
.Select(c =>
{
var isOutgoing = c.UserId == request.UserId;
var otherUserId = isOutgoing ? c.ContactId : c.UserId;
// If replica is missing, we try to at least return the record (Visibility fix part 1)
// We will handle replica creation in SendContactRequest proactively.
if (!replicas.TryGetValue(otherUserId, out var user))
{
return new ContactRequestDto(
c.Id,
new ContactUserDto(otherUserId, "Unknown", "Unknown", ""),
c.CreatedAt,
isOutgoing
);
}
return new ContactRequestDto(
c.Id,
new ContactUserDto(user.Id, user.Username, user.DisplayName, user.Avatar),
c.CreatedAt,
isOutgoing
);
})
.ToList();
return Result.Success(result);
}
}

View File

@@ -23,34 +23,57 @@ internal sealed class GetContactsQueryHandler : IQueryHandler<GetContactsQuery,
public async Task<Result<List<ContactDto>>> Handle(GetContactsQuery request, CancellationToken cancellationToken)
{
// 1. Fetch ALL accepted relations for current user
var relations = await _context.Contacts
.Where(c => (c.UserId == request.UserId || c.ContactId == request.UserId) && c.Status == ContactStatus.Accepted)
.ToListAsync(cancellationToken);
var contactIds = relations.Select(c => c.UserId == request.UserId ? c.ContactId : c.UserId).ToList();
// 2. Fetch all unique IDs for users we need replicas for
var contactIds = relations.Select(c => c.UserId == request.UserId ? c.ContactId : c.UserId).Distinct().ToList();
// 3. Fetch replicas
var replicas = await _context.UserReplicas
.Where(r => contactIds.Contains(r.Id))
.ToListAsync(cancellationToken);
.ToDictionaryAsync(r => r.Id, cancellationToken);
var result = new List<ContactDto>();
foreach (var replica in replicas)
// 4. Iterate over RELATIONS (to ensure we don't skip people with missing replicas)
foreach (var rel in relations)
{
var rel = relations.First(c => c.UserId == replica.Id || c.ContactId == replica.Id);
var otherUserId = rel.UserId == request.UserId ? rel.ContactId : rel.UserId;
result.Add(new ContactDto(
replica.Id,
replica.Username,
replica.DisplayName,
replica.Avatar,
false,
null,
rel.Id,
rel.Status == ContactStatus.Blocked,
replica.IsExternal,
replica.Domain
));
if (replicas.TryGetValue(otherUserId, out var replica))
{
result.Add(new ContactDto(
replica.Id,
replica.Username,
replica.DisplayName,
replica.Avatar,
false, // isOnline - current user query doesn't handle this here
null, // lastSeen
rel.Id,
rel.Status == ContactStatus.Blocked,
replica.IsExternal,
replica.Domain
));
}
else
{
// Return placeholder but ensure it's in the list
result.Add(new ContactDto(
otherUserId,
"Unknown",
"Unknown",
"",
false,
null,
rel.Id,
false,
false,
null
));
}
}
return Result.Success(result);

View File

@@ -1,53 +0,0 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Relations.Domain;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
using Knot.Modules.Relations.Application.Abstractions;
namespace Knot.Modules.Relations.Application.Contacts;
public record ContactUserDto(Guid Id, string Username, string DisplayName, string Avatar);
public record ContactRequestDto(Guid Id, ContactUserDto User, DateTime CreatedAt);
public record GetIncomingRequestsQuery(Guid UserId) : IQuery<List<ContactRequestDto>>;
internal sealed class GetIncomingRequestsQueryHandler : IQueryHandler<GetIncomingRequestsQuery, List<ContactRequestDto>>
{
private readonly IContactsDbContext _context;
public GetIncomingRequestsQueryHandler(IContactsDbContext context)
{
_context = context;
}
public async Task<Result<List<ContactRequestDto>>> Handle(GetIncomingRequestsQuery request, CancellationToken cancellationToken)
{
var contacts = await _context.Contacts
.Where(c => c.ContactId == request.UserId && c.Status == ContactStatus.Pending)
.ToListAsync(cancellationToken);
var requesterIds = contacts.Select(c => c.UserId).ToList();
var replicas = await _context.UserReplicas
.Where(r => requesterIds.Contains(r.Id))
.ToDictionaryAsync(r => r.Id, cancellationToken);
var result = contacts
.Where(c => replicas.ContainsKey(c.UserId))
.Select(c =>
{
var user = replicas[c.UserId];
return new ContactRequestDto(
c.Id,
new ContactUserDto(user.Id, user.Username, user.DisplayName, user.Avatar),
c.CreatedAt
);
})
.ToList();
return Result.Success(result);
}
}

View File

@@ -0,0 +1,38 @@
using Knot.Shared.Kernel;
using Knot.Shared.Kernel.Events;
using Knot.Modules.Relations.Application.Abstractions;
using Knot.Modules.Relations.Domain;
using MediatR;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
namespace Knot.Modules.Relations.Application.Contacts.Integration;
internal sealed class UserRegisteredHandler : INotificationHandler<UserRegisteredDomainEvent>
{
private readonly IContactsDbContext _context;
public UserRegisteredHandler(IContactsDbContext context)
{
_context = context;
}
public async Task Handle(UserRegisteredDomainEvent notification, CancellationToken cancellationToken)
{
var existing = await _context.UserReplicas.AnyAsync(r => r.Id == notification.UserId, cancellationToken);
if (existing) return;
var replica = UserReplica.Create(
notification.UserId,
notification.Username,
notification.DisplayName,
string.Empty, // Avatar will be synced on update
false,
null
);
_context.UserReplicas.Add(replica);
await _context.SaveChangesAsync(cancellationToken);
}
}

View File

@@ -1,4 +1,5 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Relations.Domain;
@@ -13,10 +14,12 @@ public record SendContactRequestCommand(Guid UserId, Guid ContactId) : ICommand;
internal sealed class SendContactRequestCommandHandler : ICommandHandler<SendContactRequestCommand>
{
private readonly IContactsDbContext _context;
private readonly Knot.Contracts.Auth.Infrastructure.Persistence.IAuthDbContext _authContext;
public SendContactRequestCommandHandler(IContactsDbContext context)
public SendContactRequestCommandHandler(IContactsDbContext context, Knot.Contracts.Auth.Infrastructure.Persistence.IAuthDbContext authContext)
{
_context = context;
_authContext = authContext;
}
public async Task<Result> Handle(SendContactRequestCommand request, CancellationToken cancellationToken)
@@ -26,6 +29,10 @@ internal sealed class SendContactRequestCommandHandler : ICommandHandler<SendCon
return Result.Failure(new Error("Contacts.Self", "You cannot add yourself to contacts."));
}
// 1. Proactively ensure replicas exist for both ends
await EnsureUserReplicaExists(request.UserId, cancellationToken);
await EnsureUserReplicaExists(request.ContactId, cancellationToken);
var existing = await _context.Contacts
.FirstOrDefaultAsync(f => (f.UserId == request.UserId && f.ContactId == request.ContactId) ||
(f.UserId == request.ContactId && f.ContactId == request.UserId), cancellationToken);
@@ -41,4 +48,25 @@ internal sealed class SendContactRequestCommandHandler : ICommandHandler<SendCon
return Result.Success();
}
private async Task EnsureUserReplicaExists(Guid userId, CancellationToken ct)
{
var existing = await _context.UserReplicas.AnyAsync(r => r.Id == userId, ct);
if (existing) return;
var user = await _authContext.Users.FirstOrDefaultAsync(u => u.Id == userId, ct);
if (user == null) return; // User might be external or not found
var replica = UserReplica.Create(
user.Id,
user.Username,
user.DisplayName,
user.Avatar ?? string.Empty,
user.IsExternal,
user.Domain
);
_context.UserReplicas.Add(replica);
await _context.SaveChangesAsync(ct);
}
}

View File

@@ -0,0 +1,53 @@
using System;
using System.Collections.Generic;
using System.Linq;
using Knot.Shared.Kernel;
using Microsoft.EntityFrameworkCore;
using System.Threading;
using System.Threading.Tasks;
using Knot.Modules.Relations.Application.Abstractions;
using Knot.Modules.Relations.Domain;
namespace Knot.Modules.Relations.Application.Contacts;
public record SyncReplicasCommand() : ICommand;
internal sealed class SyncReplicasCommandHandler : ICommandHandler<SyncReplicasCommand>
{
private readonly IContactsDbContext _context;
private readonly Knot.Contracts.Auth.Infrastructure.Persistence.IAuthDbContext _authContext;
public SyncReplicasCommandHandler(IContactsDbContext context, Knot.Contracts.Auth.Infrastructure.Persistence.IAuthDbContext authContext)
{
_context = context;
_authContext = authContext;
}
public async Task<Result> Handle(SyncReplicasCommand request, CancellationToken cancellationToken)
{
// 1. Fetch ALL users from Auth (since it's a repair operation)
var users = await _authContext.Users.ToListAsync(cancellationToken);
// 2. Fetch existing IDs in Replicas
var existingIds = await _context.UserReplicas.Select(r => r.Id).ToListAsync(cancellationToken);
// 3. Find missing ones
var missing = users.Where(u => !existingIds.Contains(u.Id)).ToList();
foreach (var user in missing)
{
var replica = UserReplica.Create(
user.Id,
user.Username,
user.DisplayName,
user.Avatar ?? string.Empty,
user.IsExternal,
user.Domain
);
_context.UserReplicas.Add(replica);
}
await _context.SaveChangesAsync(cancellationToken);
return Result.Success();
}
}

View File

@@ -15,6 +15,7 @@ public static class DependencyInjection
services.AddDbContext<RelationsDbContext>(options =>
options.UseNpgsql(connectionString));
services.AddScoped<IContactsDbContext>(sp => sp.GetRequiredService<RelationsDbContext>());
services.AddScoped<Knot.Contracts.Relations.Application.Abstractions.IFriendshipRepository, FriendshipRepository>();
services.AddMediatR(config =>

View File

@@ -7,10 +7,11 @@ using Knot.Shared.Kernel;
using MediatR;
using Microsoft.EntityFrameworkCore;
using Microsoft.EntityFrameworkCore.Diagnostics;
using Knot.Modules.Relations.Application.Abstractions;
namespace Knot.Modules.Relations.Infrastructure.Persistence;
public sealed class RelationsDbContext : DbContext
public sealed class RelationsDbContext : DbContext, IContactsDbContext
{
private readonly IMediator _mediator;

View File

@@ -8,6 +8,7 @@
<ItemGroup>
<ProjectReference Include="..\..\Contracts\Relations\Knot.Contracts.Relations.csproj" />
<ProjectReference Include="..\..\Contracts\Auth\Knot.Contracts.Auth.csproj" />
<ProjectReference Include="..\..\Shared\Knot.Shared.Kernel\Knot.Shared.Kernel.csproj" />
<ProjectReference Include="..\..\Shared\Knot.Shared.Infrastructure\Knot.Shared.Infrastructure.csproj" />
</ItemGroup>

View File

@@ -25,7 +25,7 @@ public static class ContactsEndpoints
group.MapGet("requests", async (ISender sender, IUserContext userContext, CancellationToken ct) =>
{
var result = await sender.Send(new GetIncomingRequestsQuery(userContext.UserId), ct);
var result = await sender.Send(new GetContactRequestsQuery(userContext.UserId), ct);
return result.IsSuccess ? Results.Ok(result.Value) : Results.BadRequest(result.Error);
});
@@ -86,6 +86,7 @@ public static class ContactsEndpoints
return result.IsSuccess ? Results.Ok(new { success = true }) : Results.NotFound(result.Error.Description);
});
group.MapDelete("{id:guid}", async ([FromRoute] Guid id, ISender sender, IUserContext userContext, CancellationToken ct) =>
{
var result = await sender.Send(new RemoveContactCommand(userContext.UserId, id), ct);