using gehGassiApp.Core.Interfaces; using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; using gehGassiApp.Core.Data; using gehGassiApp.Core.Interfaces.Synchronization; using gehGassiApp.Domain.Messages; using gehGassiApp.Core.Services.Synchronization; using gehGassiApp.Domain.Common; using gehGassiApp.Domain.News; using CommunityToolkit.Mvvm.Messaging; using gehGassiApp.Domain.Users; using System.Text.Json; using gehGassi.Dto.Messages; using gehGassiApp.Core.Helper; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Metadata.Internal; namespace gehGassiApp.Core.Services { /// /// Schnittstellenbeschreibung für einen Service der die verwaltung von Nachrichten ermöglicht /// public class MessageService : ServiceBase, IMessageService { private readonly ICommunicationService _communicationService; private readonly ISyncInfoPullService _pullService; private readonly ISyncInfoPushService _pushService; private readonly IRepository _appUserRepository; private readonly IRepository _conversationRepository; /// /// Erstellt eine Instanz /// /// Instanz eines IUnitOfWork /// Instanz eines ICommunicationService /// Instanz eines ISyncInfoPullService /// Instanz eines ISyncInfoPushService public MessageService(IUnitOfWork unitOfWork, ICommunicationService communicationService, ISyncInfoPullService syncInfoPullService, ISyncInfoPushService syncInfoPushService) : base(unitOfWork) { _communicationService = communicationService; _pullService = syncInfoPullService; _pushService = syncInfoPushService; _appUserRepository = unitOfWork.GetRepository(); _conversationRepository = unitOfWork.GetRepository(); } /// /// Gibt eine Entität anhand der eindeutigen Id zurück /// /// Id der Entität /// Gibt an ob NoTRacking verwendet werden soll. Es werden keine Entitäten im EF-Speicher gehalten /// Entität oder null, wenn nicht gefunden public override Message Get(object id, bool noTracking = true) { return Repository.SingleOrDefault(c => c.Id == id.ToString(), noTracking); } /// /// Gibt eine Entität anhand der eindeutigen Id zurück /// /// Id der Entität /// Gibt an ob NoTRacking verwendet werden soll. Es werden keine Entitäten im EF-Speicher gehalten /// Entität oder null, wenn nicht gefunden public override async Task GetAsync(object id, bool noTracking = true) { return await Repository.SingleOrDefaultAsync(c => c.Id == id.ToString(), noTracking).ConfigureAwait(false); } /// /// Gibt eine Liste der Nachrichten für eine Konversation zurück /// /// Id der Konversation /// Aktuelles Accesstoken /// CancellationToken /// Id des Senders /// Liste Nachrichten public async Task> GetMessagesAsync(string senderId, string conversationId, string accessToken, CancellationToken token) { var isConnected = await _communicationService.IsConnected(); //isConnected = false; if (isConnected) { var messagesResult = await _communicationService.GetMessagesAsync(senderId, conversationId, null, accessToken, token); if (messagesResult.Success) { var result = await HandleChangesAsync(messagesResult.Value); if (result.LastUpdate != null) await _communicationService.ConfirmMessagesAsync(senderId, result.LastUpdate, accessToken, token); } } //Jetzt die DB abfragen var messages = (await Repository.FindAsync(c => c.ConversationId == conversationId && c.Deleted == false, noTracking: true).ConfigureAwait(false)).ToList(); return messages; } /// /// Gibt eine Liste der Nachrichten für eine Konversation zurück /// /// Id der Konversation /// Wie viele News sollen abgerufen werden? -1 Wenn nicht anwenden. /// Wie viele News sollen ausgelassen werden? -1 Wenn nicht anwenden. /// Aktuelles Accesstoken /// CancellationToken /// Id des Senders /// Liste Nachrichten public async Task>> GetMessagesAsync(string senderId, string conversationId, int take, int skip, string accessToken, CancellationToken token) { var result = new ListCommunicationResult>() { Value = new List() }; var isConnected = await _communicationService.IsConnected(); //isConnected = false; if (isConnected) { var messagesResult = await _communicationService.GetMessagesAsync(senderId, conversationId, null, accessToken, token).ConfigureAwait(false); if (messagesResult.Success) { var changesResult = await HandleChangesAsync(messagesResult.Value); if (changesResult.LastUpdate != null) await _communicationService.ConfirmMessagesAsync(senderId, changesResult.LastUpdate, accessToken, token); } } //Jetzt die DB abfragen result.Total = await Repository.CountAsync(c => c.ConversationId == conversationId && c.Deleted == false); //var query = Repository.Query(c => c.ConversationId == conversationId && c.Deleted == false).OrderByDescending(c => c.UpdatedAt); //var messages = await (query.Skip(skip).Take(take).ToListAsync().ConfigureAwait(false)); var messages = (await Repository.FindAsync(c => c.ConversationId == conversationId && c.Deleted == false, c => c.UpdatedAt, skip, take, false)).ToList(); result.Success = true; result.Value = messages; result.Take = take; result.Skip = skip; return result; } /// /// Gibt eine Liste der Nachrichten für die Synchronisation zurück /// /// Aktuelles Accesstoken /// CancellationToken /// Id des Senders /// Liste Nachrichten public async Task>> GetMessagesForSyncAsync(string senderId, string accessToken, CancellationToken token) { var result = new CommunicationResult>() { Value = new List() }; var isConnected = await _communicationService.IsConnected(); //isConnected = false; if (isConnected) { var messagesResult = await _communicationService.GetMessagesAsync(senderId, null, accessToken, token); if (messagesResult.Success) { result.Value = messagesResult.Value; result.Success = true; } } return result; } /// /// Erstellen einer Nachricht. Wird lokal gespeichert und an der Server übertragen. /// Wenn keine Verbindung besteht, wird ein Delta erstellt /// /// id des Senders /// Id des Empfängers /// Id der Konversation /// Inhalt /// Aktuelles Accesstoken /// CancellationToken /// Task public async Task AddMessageAsync(string senderId, string receiverId, string conversationId, string content, string accessToken, CancellationToken token) { var messageId = Guid.NewGuid().ToString("N"); var message = new Message() { Id = messageId, UpdatedAt = DateTimeOffset.UtcNow, Direction = MessageDirection.Out, Read = true, Content = content, ConversationId = conversationId, Deleted = false }; Repository.Add(message); await CommitAsync(); //Jetzt die Nachricht an der Server übertragen oder ein Delta erstellen var isConnected = await _communicationService.IsConnected(); var createDelta = true; if (isConnected) { var createResult = await _communicationService.AddMessageAsync(senderId, receiverId, message, accessToken, token); if (createResult.Success) createDelta = false; } if (createDelta) { await _pushService.AddAsync(nameof(Message), message.Id, SyncOperation.Create, message).ConfigureAwait(false); } return message; } /// /// Setzt alle Nachrichten in einer Konversation als gelesen /// /// Id der Konversation /// Task public async Task SetMessagesReadAsync(string conversationId) { var changedCount = await Repository.ExecuteRawSqlAsync("UPDATE Messages SET Read = 1 WHERE ConversationId = {0} AND Read = 0", new object[]{conversationId} ); await CommitAsync(); var changedConversationCount = await _conversationRepository.ExecuteRawSqlAsync("UPDATE Conversations SET HasUnreadMessages = 0 WHERE Id = {0}", new object[] { conversationId }); await CommitAsync(); } /// /// Gibt eine Liste ungelesener Nachrichten zurück /// /// Liste ungelesener Nachrichten public async Task> GetUnreadAsync() { var unreadMessages = await Repository.FindAsync(c => c.Read == false).ConfigureAwait(false); return unreadMessages.ToList(); } /// /// Gibt eine Nachricht einer Konversation zurück /// /// Id der Konversation /// Id der Nachricht /// Nachricht oder null, wenn nicht gefunden public async Task GetAsync(string conversationId, string messageId) { return await Repository.FirstOrDefaultAsync(c => c.ConversationId == conversationId && c.Id == messageId).ConfigureAwait(false); } /// /// Löscht eine Nachricht. /// Die Nachricht wird auf gelöscht gesetzt. /// /// Id der Konversation /// Id der Nachricht /// True wenn wirklich gelöscht werden soll /// True wenn erfolgreich, false sonst public async Task DeleteAsync(string conversationId, string messageId, bool forceDelete = false) { var message = await Repository.FirstOrDefaultAsync(c => c.ConversationId == conversationId && c.Id == messageId, false).ConfigureAwait(false); if (message != null) { if (forceDelete) { Repository.Remove(message); } else { message.Deleted = true; Repository.Update(message); } await CommitAsync(); return true; } return false; } #region Pull-Push Implementation /// /// Holen der Daten vom lokalen Speicher die noch nicht synchronisiert wurden und senden an den Server /// /// Aktuelles Accesstoken /// CancellationToken /// Task async Task> ISyncPushPullService.PushAsync(string accessToken, CancellationToken token) { var result = new SyncResult(); if (await _pushService.HasOpenAsync(nameof(Message)) > 0) { var deltas = await _pushService.GetAllAsync(nameof(Message)); if (deltas.Any()) { var user = await _appUserRepository.FirstOrDefaultAsync(c => c.Id != "", false).ConfigureAwait(false); deltas = deltas.OrderBy(c => c.DateTime).ToList(); Conversation conversation = null; var isConnected = await _communicationService.IsConnected(); if (!isConnected) return result; foreach (var syncInfoPush in deltas) { //Löschen einer Nachricht geben wir nicht weiter... Update gibt es eigentlich auch keines if (syncInfoPush.Operation == SyncOperation.Delete || syncInfoPush.Operation == SyncOperation.Edit) { await _pushService.RemoveAsync(syncInfoPush.Id).ConfigureAwait(false); continue; } try { if (!string.IsNullOrWhiteSpace(syncInfoPush.Value)) { var message = JsonSerializer.Deserialize(syncInfoPush.Value, new JsonSerializerOptions(JsonSerializerDefaults.Web)); if (conversation == null || conversation.Id != message.ConversationId) conversation = await _conversationRepository.GetAsync(message.ConversationId); var createResult = await _communicationService.AddMessageAsync(user.Id, conversation.Recipient, message, accessToken, token); if (createResult.Success) { await _pushService.RemoveAsync(syncInfoPush.Id).ConfigureAwait(false); } else { //Fehler, weg damit await _pushService.RemoveAsync(syncInfoPush.Id).ConfigureAwait(false); } } else { await _pushService.RemoveAsync(syncInfoPush.Id).ConfigureAwait(false); } } catch (Exception ex) { System.Diagnostics.Debug.WriteLine(ex.Message); await _pushService.RemoveAsync(syncInfoPush.Id).ConfigureAwait(false); } } //Else ist nichts tun, konnte nicht übertragen werden. } } return result; } /// /// Holen der letzten Daten vom Server und Synchronisieren mit den lokalen Daten /// /// Sprache /// Aktuelles Accesstoken /// CancellationToken /// Task async Task> ISyncPushPullService.PullAsync(string language, string accessToken, CancellationToken token) { var result = new SyncResult(); var isConnected = await _communicationService.IsConnected(); if (isConnected) { //Zuerst lezte Aktivität holen... HIER NICHT da wir immer alle offenen Nachrichten holen //DateTimeOffset? lastUpdate = null; //var lastSyncInfo = await _pullService.GetAsync(nameof(Message)); //if (lastSyncInfo != null) // lastUpdate = lastSyncInfo.LastUpdate; var user = await _appUserRepository.FirstOrDefaultAsync(c => c.Id != "", false).ConfigureAwait(false); var messagesResult = await _communicationService.GetMessagesAsync(user.Id, null, accessToken, token); if (messagesResult.Success) { result = await HandleChangesAsync(messagesResult.Value); if (result.LastUpdate != null) { await _communicationService.ConfirmMessagesAsync(user.Id, result.LastUpdate, accessToken, token); } if (result.Added.Any()) { var groups = result.Added.GroupBy(c => c.ConversationId); foreach (var group in groups) { var conversationId = group.Key; var lastMessage = group.MaxBy(c => c.UpdatedAt); if (lastMessage != null) { var conversation = await _conversationRepository.GetAsync(conversationId); if (conversation != null) { conversation.LastMessageId = lastMessage.Id; conversation.LastMessagePreview = lastMessage.Content.Ellipsis(100); conversation.HasUnreadMessages = true; _conversationRepository.Update(conversation); } } } await CommitAsync().ConfigureAwait(false); } } //result.LastUpdate ??= lastUpdate; } return result; } #endregion #region Private /// /// Behandeln der Liste von Nachrichten wenn welche vom Online-Store geholt werden. /// /// Liste der Konversationen /// Task private async Task> HandleChangesAsync(List messages) { var result = new SyncResult(); foreach (var message in messages) { var localMessage = await Repository.FirstOrDefaultAsync(c => c.Id == message.Id, false).ConfigureAwait(false); if (localMessage != null) { if (result.LastUpdate == null || result.LastUpdate < localMessage.UpdatedAt) result.LastUpdate = localMessage.UpdatedAt; result.Updated.Add(localMessage); } if (localMessage == null) { //TODO: Prüfen ob die Konversation noch existiert...?? var conversation = await _conversationRepository.FirstOrDefaultAsync(c => c.Id == message.ConversationId); if (conversation != null && !conversation.Deleted) { message.Read = false; Repository.Add(message); if (result.LastUpdate == null || result.LastUpdate < message.UpdatedAt) result.LastUpdate = message.UpdatedAt; result.Added.Add(message); } } } if (result.HasChanges) { try { await CommitAsync().ConfigureAwait(false); } catch (Exception ex) { var err = ex.Message; } } return result; } #endregion } }