using System.Text.Json; using AdminBLL.Subscription; using AdminDAL.DTO.Subscription; using ECPBLL.Communication.Interfaces; using ECPDAL.Communication.DTOs; using ECPDAL.Communication.Interfaces; using GB5Shared.PushNotification; using GB5Shared.DTO.PushNotification; using GB5Shared.DaprCache; using GB5Shared.DTO.Framework.Login; using GB5Shared.EventLogPublish; using GB5Shared.GB5Constant; using GB5Shared.GenerateAutoNumber; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using GB5Shared.Validation; using Microsoft.Extensions.Logging; namespace ECPBLL.Communication.Implementations; public class CommentThreadBLL( ICommentThreadDAL _dal, ICommentMentionDAL _mentionDal, IValidation _validation, AutoNumber _autoNumber, KeyInvalidate _keyInvalidate, EventLogPublish _eventLog, ISubscriptionBLL _subscriptionBll, IPushNotificationBLL _pushNotificationBll, ICommunicationPushService _pushService, ILogger _logger ) : ICommentThreadBLL { // PropertyNamingPolicy left at its default (preserve declared PascalCase names) -- GB5's // wire-format convention (see CLAUDE.md "PropertyNamingPolicy = null") is PascalCase // everywhere; this used to override to camelCase, silently breaking every frontend consumer. private static readonly JsonSerializerOptions JsonOpts = new(); // ── GET ─────────────────────────────────────────────────────────── public async Task GetCommentThreadAsync(int commentThreadId, LoginDTO login, CancellationToken ct) { GB5Trace.Step("get-comment-thread", new { commentThreadId }); var thread = await _dal.GetCommentThreadAsync(commentThreadId, login, ct).ConfigureAwait(false); if (thread is null) throw new InvalidOperationException($"Comment thread {commentThreadId} not found."); if (!await CanViewCommentThreadAsync(thread, login, ct).ConfigureAwait(false)) throw new UnauthorizedAccessException("You do not have access to this comment thread."); return JsonSerializer.Serialize(thread, JsonOpts); } public async Task GetCommentThreadTreeAsync(int objectTypeId, int objectId, LoginDTO login, CancellationToken ct) { GB5Trace.Step("get-comment-thread-tree", new { objectTypeId, objectId }); var rootThreads = await _dal.GetCommentThreadListByObjectAsync(objectTypeId, objectId, login, ct) .ConfigureAwait(false); var allTrees = new List(); foreach (var root in rootThreads) { var flat = await _dal.GetCommentThreadTreeAsync(root.RootCommentThreadId, login, ct).ConfigureAwait(false); var visible = new List(); foreach (var row in flat) { if (await CanViewCommentThreadAsync(row, login, ct).ConfigureAwait(false)) visible.Add(row); } allTrees.AddRange(BuildTree(visible)); } return JsonSerializer.Serialize(allTrees, JsonOpts); } public async Task GetCommentThreadListAsync(int objectTypeId, int objectId, LoginDTO login, CancellationToken ct) { GB5Trace.Step("get-comment-thread-list", new { objectTypeId, objectId }); var roots = await _dal.GetCommentThreadListByObjectAsync(objectTypeId, objectId, login, ct).ConfigureAwait(false); var visible = new List(); foreach (var root in roots) { if (await CanViewCommentThreadAsync(root, login, ct).ConfigureAwait(false)) visible.Add(root); } return JsonSerializer.Serialize(visible, JsonOpts); } /// Assembles a flat, authorized row list into a parent/child tree in memory — /// avoids any recursive-CTE cross-engine SQL (see CommentThreadQB.GET_COMMENT_THREAD_TREE). /// Rows whose parent was filtered out by CanView (or never existed) surface as roots too, so /// an authorized reply is never silently dropped just because its parent isn't visible. private static List BuildTree(List rows) { var byId = rows.ToDictionary(r => r.CommentThreadId); var roots = new List(); foreach (var row in rows) row.Replies = []; foreach (var row in rows) { if (row.ParentCommentThreadId != -1 && byId.TryGetValue(row.ParentCommentThreadId, out var parent)) parent.Replies.Add(row); else roots.Add(row); } return roots; } // ── SAVE ────────────────────────────────────────────────────────── public async Task SaveCommentThreadAsync(SaveCommentThreadDTO dto, LoginDTO login, CancellationToken ct) { GB5Trace.Step("validate-comment-thread", new { dto.ParentCommentThreadId, dto.ObjectTypeId, dto.ObjectId }); await _validation.NotEmpty(dto.CommentText, nameof(dto.CommentText)).ConfigureAwait(false); if (dto.ObjectTypeId == 0) throw new InvalidOperationException("ObjectTypeId is required."); if (dto.ObjectId == 0) throw new InvalidOperationException("ObjectId is required."); try { var isReply = dto.ParentCommentThreadId != -1; var rootCommentThreadId = 0; if (isReply) { var parent = await _dal.GetCommentThreadAsync(dto.ParentCommentThreadId, login, ct).ConfigureAwait(false); if (parent is null) throw new InvalidOperationException($"Parent comment thread {dto.ParentCommentThreadId} not found."); if (!await CanViewCommentThreadAsync(parent, login, ct).ConfigureAwait(false)) throw new UnauthorizedAccessException("You do not have access to this comment thread."); rootCommentThreadId = parent.RootCommentThreadId; } var autoNumber = await _autoNumber.GetNumberAsync(1, "COMMENTTHREAD", login).ConfigureAwait(false); var newId = autoNumber.StartNumber; var entity = new CommentThreadDTO { CommentThreadId = newId, ParentCommentThreadId = isReply ? dto.ParentCommentThreadId : -1, RootCommentThreadId = isReply ? rootCommentThreadId : newId, ObjectTypeId = dto.ObjectTypeId, ObjectId = dto.ObjectId, UserId = login.UserId, UserCode = login.UserCode, UserName = login.UserName, CommentText = dto.CommentText, CommentStatus = CommunicationConstants.CommentStatus.Open, AccessType = dto.AccessType == 0 ? CommunicationConstants.AccessType.Public : dto.AccessType, PrivateUserGroupId = dto.PrivateUserGroupId, SecurityMarkId = dto.SecurityMarkId, BizTransactionId = dto.BizTransactionId, OuId = dto.OuId, TenantId = login.ClientId, DatabaseName = login.DatabaseName, DatabaseType = login.DatabaseType, Version = 0, Status = CommunicationConstants.RowStatus.Active, CreatedById = login.UserId, CreatedOn = DateTime.UtcNow, ModifiedById = login.UserId, ModifiedOn = DateTime.UtcNow, }; GB5Trace.Step("save-comment-thread", new { newId, isReply }); await _dal.SaveCommentThreadAsync(entity, login, ct).ConfigureAwait(false); var newMentions = await SaveMentionsAsync(entity, dto.MentionedUserIds, login, ct).ConfigureAwait(false); var eventTypeId = isReply ? Constant.EventTypeConstant.COMMENTTHREADREPLIEDEVENTTYPEID : Constant.EventTypeConstant.COMMENTTHREADCREATEDEVENTTYPEID; GB5Trace.Step("event-publish", new { eventTypeId }); await _eventLog.PublishEventLogAsync( isReply ? "Comment Thread Replied" : "Comment Thread Created", entity, eventTypeId, newId, login, ct: ct).ConfigureAwait(false); await InvalidateThreadCacheAsync(dto.ObjectTypeId, dto.ObjectId, login).ConfigureAwait(false); if (isReply) await _pushService.PushReplyAsync(dto.ObjectTypeId, dto.ObjectId, entity, ct).ConfigureAwait(false); else await _pushService.PushCommentAsync(dto.ObjectTypeId, dto.ObjectId, entity, ct).ConfigureAwait(false); await NotifyAudienceAsync(entity, newMentions, login, ct).ConfigureAwait(false); return $"{SuccessResponse.SaveSuccessMessage} {newId}"; } catch (Exception ex) { GB5Trace.MarkFailed("save-comment-thread-failed", ex); _logger.LogError(ex, "SaveCommentThread failed for ObjectType {ObjectTypeId} Object {ObjectId}", dto.ObjectTypeId, dto.ObjectId); throw; } } private async Task> SaveMentionsAsync( CommentThreadDTO thread, List mentionedUserIds, LoginDTO login, CancellationToken ct) { var saved = new List(); foreach (var mentionedUserId in mentionedUserIds.Distinct()) { if (mentionedUserId == 0 || mentionedUserId == login.UserId) continue; var mentionAutoNumber = await _autoNumber.GetNumberAsync(1, "COMMENTMENTION", login).ConfigureAwait(false); var mention = new CommentMentionDTO { CommentMentionId = mentionAutoNumber.StartNumber, CommentThreadId = thread.CommentThreadId, MentionedUserId = mentionedUserId, IsNotified = false, TenantId = login.ClientId, DatabaseName = login.DatabaseName, DatabaseType = login.DatabaseType, Status = CommunicationConstants.RowStatus.Active, CreatedById = login.UserId, CreatedOn = DateTime.UtcNow, ModifiedById = login.UserId, ModifiedOn = DateTime.UtcNow, ObjectTypeId = thread.ObjectTypeId, ObjectId = thread.ObjectId, CommentText = thread.CommentText, }; await _mentionDal.SaveCommentMentionAsync(mention, login, ct).ConfigureAwait(false); saved.Add(mention); // GetPendingCommentMentions is cached USER_LEVEL per mentioned user — without this, // a user's mention inbox (and the topbar bell badge count) would stay stale until // cache TTL expiry after every new @mention, regardless of who created the comment. await InvalidatePendingMentionsCacheAsync(mentionedUserId, login).ConfigureAwait(false); } return saved; } // GetPendingCommentMentions' own cache key is built from the CALLER's LoginDTO (its UserId // becomes both the ObjectId and part of the USER_LEVEL hash) — so invalidating the mentioned // user's cache requires a LoginDTO shaped as if THEY were logged in, not the comment author. // Every other field CacheKeyGeneration reads for USER_LEVEL (ServerId/ClientId/ // ServerConfigId/RoleId/LanguageId) is copied from the real login since those describe the // tenant/session context, not "who is asking" — only UserId legitimately differs here. private async Task InvalidatePendingMentionsCacheAsync(int mentionedUserId, LoginDTO login) { var mentionedUserLogin = new LoginDTO { ServerId = login.ServerId, ClientId = login.ClientId, ServerConfigId = login.ServerConfigId, RoleId = login.RoleId, LanguageId = login.LanguageId, UserId = mentionedUserId, }; var cacheKey = new CacheKeyGeneration().KeyGeneration( mentionedUserId, Constant.EntityConstant.OBJECTPENDINGCOMMENTMENTION, Constant.CacheKeyLevel.USER_LEVEL, mentionedUserLogin); await _keyInvalidate.InvalidateCache(cacheKey).ConfigureAwait(false); } // ── UPDATE STATUS ───────────────────────────────────────────────── public async Task UpdateCommentThreadStatusAsync(UpdateCommentThreadStatusDTO dto, LoginDTO login, CancellationToken ct) { GB5Trace.Step("validate-comment-thread-status", new { dto.CommentThreadId, dto.NewStatus }); if (dto.NewStatus is < 1 or > 3) throw new InvalidOperationException("NewStatus must be 1 (open), 2 (resolved), or 3 (archive)."); try { var thread = await _dal.GetCommentThreadAsync(dto.CommentThreadId, login, ct).ConfigureAwait(false); if (thread is null) throw new InvalidOperationException($"Comment thread {dto.CommentThreadId} not found."); if (!await CanViewCommentThreadAsync(thread, login, ct).ConfigureAwait(false)) throw new UnauthorizedAccessException("You do not have access to this comment thread."); GB5Trace.Step("save-comment-thread-status", new { dto.CommentThreadId, dto.NewStatus }); await _dal.UpdateCommentThreadStatusAsync(dto.CommentThreadId, dto.NewStatus, login.UserId, login, ct) .ConfigureAwait(false); GB5Trace.Step("event-publish", new { Constant.EventTypeConstant.COMMENTTHREADSTATUSEVENTTYPEID }); await _eventLog.PublishEventLogAsync( "Comment Thread Status Changed", dto, Constant.EventTypeConstant.COMMENTTHREADSTATUSEVENTTYPEID, dto.CommentThreadId, login, ct: ct) .ConfigureAwait(false); await InvalidateThreadCacheAsync(thread.ObjectTypeId, thread.ObjectId, login).ConfigureAwait(false); await InvalidateCommentDetailCacheAsync(dto.CommentThreadId, login).ConfigureAwait(false); await _pushService.PushStatusChangeAsync(thread.ObjectTypeId, thread.ObjectId, dto.CommentThreadId, dto.NewStatus, ct) .ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { GB5Trace.MarkFailed("update-comment-thread-status-failed", ex); _logger.LogError(ex, "UpdateCommentThreadStatus failed for CommentThreadId {CommentThreadId}", dto.CommentThreadId); throw; } } // ── DELETE ──────────────────────────────────────────────────────── public async Task DeleteCommentThreadAsync(int commentThreadId, LoginDTO login, CancellationToken ct) { GB5Trace.Step("delete-comment-thread", new { commentThreadId }); try { var thread = await _dal.GetCommentThreadAsync(commentThreadId, login, ct).ConfigureAwait(false); if (thread is null) throw new InvalidOperationException($"Comment thread {commentThreadId} not found."); // Phase 1 rule: only the original author may delete their own comment thread. (Not // spelled out explicitly by the plan; a documented simplification rather than // reusing CanView, since "can see it" and "can delete it" are different rights.) if (thread.UserId != login.UserId) throw new UnauthorizedAccessException("Only the author can delete this comment thread."); await _dal.DeleteCommentThreadAsync(commentThreadId, login.UserId, login, ct).ConfigureAwait(false); GB5Trace.Step("event-publish", new { Constant.EventTypeConstant.COMMENTTHREADDELETEDEVENTTYPEID }); await _eventLog.PublishEventLogAsync( "Comment Thread Deleted", new { commentThreadId }, Constant.EventTypeConstant.COMMENTTHREADDELETEDEVENTTYPEID, commentThreadId, login, ct: ct) .ConfigureAwait(false); await InvalidateThreadCacheAsync(thread.ObjectTypeId, thread.ObjectId, login).ConfigureAwait(false); await InvalidateCommentDetailCacheAsync(commentThreadId, login).ConfigureAwait(false); // Reuses the status-change push (no dedicated "deleted" hub method in the plan's // client contract) — NewStatus=0 is a client-side sentinel meaning "removed", never // persisted (COMMENTSTATUS itself keeps its last real value in the DB). await _pushService.PushStatusChangeAsync(thread.ObjectTypeId, thread.ObjectId, commentThreadId, 0, ct) .ConfigureAwait(false); return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { GB5Trace.MarkFailed("delete-comment-thread-failed", ex); _logger.LogError(ex, "DeleteCommentThread failed for CommentThreadId {CommentThreadId}", commentThreadId); throw; } } // ── VISIBILITY (plan: "Access, security level & notification routing") ─────────────────── public async Task CanViewCommentThreadAsync(CommentThreadDTO thread, LoginDTO login, CancellationToken ct) { // 1. Thread level — most specific, checked first. if (thread.AccessType == CommunicationConstants.AccessType.Private && thread.PrivateUserGroupId != -1) { var inGroup = await _dal.IsUserInGroupAsync(login.UserId, thread.PrivateUserGroupId, login, ct) .ConfigureAwait(false); if (!inGroup && thread.UserId != login.UserId) return false; } if (thread.SecurityMarkId != -1) { var hasRights = await _dal.HasSecurityRightsAsync( login.UserId, login.RoleId, CommunicationConstants.CommentThreadEntityId, thread.SecurityMarkId, login, ct) .ConfigureAwait(false); if (!hasRights && thread.UserId != login.UserId) return false; } // 2. BizTransaction level — deferred to the owning module per the plan (Communication // does not re-implement business-transaction rights); no check performed here. // 3. OU level — viewer's WorkOUId must be the thread's OU or one of its descendants, // resolved via a bounded walk over MORGANIZATIONGROUPDETAIL.ORGANISATIONGROUPPARENTID. if (thread.OuId != -1) { var inScope = await IsOuInScopeAsync(thread.OuId, login.WorkOUId, login, ct).ConfigureAwait(false); if (!inScope) return false; } // 4. ObjectType/Object floor — delegated to the owning module per the plan (Communication // doesn't own business data); not re-implemented here. return true; } /// Bounded C# walk over MORGANIZATIONGROUPDETAIL (rather than a recursive CTE, to /// avoid PostgreSQL/SQL Server WITH RECURSIVE dialect differences) — resolves every OU /// reachable from threadOuId's organization group(s) by following child groups whose /// ORGANISATIONGROUPPARENTID points back up the chain. private async Task IsOuInScopeAsync(int threadOuId, int viewerOuId, LoginDTO login, CancellationToken ct) { if (viewerOuId == threadOuId) return true; var details = await _dal.GetOrganizationGroupDetailsAsync(login, ct).ConfigureAwait(false); var frontier = details .Where(d => d.OuId == threadOuId) .Select(d => d.OrganizationGroupId) .Distinct() .ToHashSet(); var visitedGroups = new HashSet(frontier); var resolvedOuIds = new HashSet { threadOuId }; for (var depth = 0; depth < CommunicationConstants.MaxOuHierarchyDepth && frontier.Count > 0; depth++) { var nextFrontier = new HashSet(); foreach (var groupId in frontier) { foreach (var child in details.Where(d => d.OrganizationGroupParentId == groupId)) { if (child.OuId != -1) resolvedOuIds.Add(child.OuId); if (visitedGroups.Add(child.OrganizationGroupId)) nextFrontier.Add(child.OrganizationGroupId); } } frontier = nextFrontier; } return resolvedOuIds.Contains(viewerOuId); } // ── NOTIFICATION AUDIENCE (plan: mentions + participants + MSUBSCRIPTION) ──────────────── private async Task NotifyAudienceAsync( CommentThreadDTO thread, List newMentions, LoginDTO login, CancellationToken ct) { try { var audience = new HashSet(); foreach (var mention in newMentions) audience.Add(mention.MentionedUserId); var participants = await _dal.GetThreadParticipantUserIdsAsync(thread.RootCommentThreadId, login, ct) .ConfigureAwait(false); foreach (var participantId in participants) audience.Add(participantId); try { var subsJson = await _subscriptionBll .GetSubscriptionsAsync(thread.ObjectTypeId, thread.ObjectId, login, ct) .ConfigureAwait(false); var subs = JsonSerializer.Deserialize>(subsJson, JsonOpts) ?? []; foreach (var sub in subs) { if (string.IsNullOrWhiteSpace(sub.SubscriptionUserIds)) continue; foreach (var idText in sub.SubscriptionUserIds.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries)) if (int.TryParse(idText, out var uid)) audience.Add(uid); } } catch (Exception ex) { _logger.LogWarning(ex, "Subscription audience resolution failed for ObjectType {ObjectTypeId} Object {ObjectId}", thread.ObjectTypeId, thread.ObjectId); } audience.Remove(login.UserId); foreach (var mention in newMentions) await _pushService.PushMentionAsync(mention.MentionedUserId, mention, ct).ConfigureAwait(false); if (newMentions.Count > 0) { await _mentionDal.MarkMentionsNotifiedAsync(thread.CommentThreadId, login, ct).ConfigureAwait(false); GB5Trace.Step("event-publish", new { Constant.EventTypeConstant.COMMENTMENTIONEDEVENTTYPEID }); await _eventLog.PublishEventLogAsync( "User Mentioned in Comment", newMentions, Constant.EventTypeConstant.COMMENTMENTIONEDEVENTTYPEID, thread.CommentThreadId, login, ct: ct) .ConfigureAwait(false); } if (audience.Count > 0) { try { await _pushNotificationBll.SendNotificationToUsers(new PushNotificationDTO { NotificationTitle = "New comment activity", NotificationBody = thread.CommentText.Length > 200 ? thread.CommentText[..200] : thread.CommentText, NotificationTopic = $"COMMUNICATION_OBJECT_{thread.ObjectTypeId}_{thread.ObjectId}", UserIds = audience.ToList(), }).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("push-notification-fallback-failed", ex); _logger.LogWarning(ex, "PushNotification fallback failed for thread {CommentThreadId}", thread.CommentThreadId); } } } catch (Exception ex) { // Notification-audience resolution must never fail the underlying save. GB5Trace.MarkFailed("notify-audience-failed", ex); _logger.LogWarning(ex, "NotifyAudience failed for thread {CommentThreadId}", thread.CommentThreadId); } } // Clears GetCommentThreadTree/GetCommentThreadList (both share OBJECTCOMMENTTHREADLIST, keyed // by "{objectTypeId}:{objectId}") — every mutation that changes what's IN the thread for this // object needs this, including the CollaborationSpaceDiscussion/MeetingMinutes wrappers that // call through to this same BLL with a fixed ObjectTypeId. private async Task InvalidateThreadCacheAsync(int objectTypeId, int objectId, LoginDTO login) { var cacheKey = new CacheKeyGeneration().KeyGeneration( $"{objectTypeId}:{objectId}", Constant.EntityConstant.OBJECTCOMMENTTHREADLIST, Constant.CacheKeyLevel.CLIENT_LEVEL, login); await _keyInvalidate.InvalidateCache(cacheKey).ConfigureAwait(false); } // Clears GetCommentThread's own per-id cache — needed by UpdateCommentThreadStatus and // DeleteCommentThread, both of which mutate a specific existing comment row directly // (SaveCommentThread never needs this: a brand-new reply has no prior single-item cache, and // editing an existing comment isn't a supported operation on this endpoint). private async Task InvalidateCommentDetailCacheAsync(int commentThreadId, LoginDTO login) { var cacheKey = new CacheKeyGeneration().KeyGeneration( commentThreadId, Constant.EntityConstant.OBJECTCOMMENTTHREAD, Constant.CacheKeyLevel.CLIENT_LEVEL, login); await _keyInvalidate.InvalidateCache(cacheKey).ConfigureAwait(false); } }