466 lines
21 KiB
C#
466 lines
21 KiB
C#
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
|
|
{
|
|
/// <summary>
|
|
/// Schnittstellenbeschreibung für einen Service der die verwaltung von Nachrichten ermöglicht
|
|
/// </summary>
|
|
public class MessageService : ServiceBase<Message>, IMessageService
|
|
{
|
|
private readonly ICommunicationService _communicationService;
|
|
private readonly ISyncInfoPullService _pullService;
|
|
private readonly ISyncInfoPushService _pushService;
|
|
private readonly IRepository<AppUser> _appUserRepository;
|
|
private readonly IRepository<Conversation> _conversationRepository;
|
|
|
|
/// <summary>
|
|
/// Erstellt eine Instanz
|
|
/// </summary>
|
|
/// <param name="unitOfWork">Instanz eines IUnitOfWork</param>
|
|
/// <param name="communicationService">Instanz eines ICommunicationService</param>
|
|
/// <param name="syncInfoPullService">Instanz eines ISyncInfoPullService</param>
|
|
/// <param name="syncInfoPushService">Instanz eines ISyncInfoPushService</param>
|
|
public MessageService(IUnitOfWork unitOfWork, ICommunicationService communicationService, ISyncInfoPullService syncInfoPullService, ISyncInfoPushService syncInfoPushService) : base(unitOfWork)
|
|
{
|
|
_communicationService = communicationService;
|
|
_pullService = syncInfoPullService;
|
|
_pushService = syncInfoPushService;
|
|
_appUserRepository = unitOfWork.GetRepository<AppUser>();
|
|
_conversationRepository = unitOfWork.GetRepository<Conversation>();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Entität anhand der eindeutigen Id zurück
|
|
/// </summary>
|
|
/// <param name="id">Id der Entität</param>
|
|
/// <param name="noTracking">Gibt an ob NoTRacking verwendet werden soll. Es werden keine Entitäten im EF-Speicher gehalten</param>
|
|
/// <returns>Entität oder null, wenn nicht gefunden</returns>
|
|
public override Message Get(object id, bool noTracking = true)
|
|
{
|
|
return Repository.SingleOrDefault(c => c.Id == id.ToString(), noTracking);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Entität anhand der eindeutigen Id zurück
|
|
/// </summary>
|
|
/// <param name="id">Id der Entität</param>
|
|
/// <param name="noTracking">Gibt an ob NoTRacking verwendet werden soll. Es werden keine Entitäten im EF-Speicher gehalten</param>
|
|
/// <returns>Entität oder null, wenn nicht gefunden</returns>
|
|
public override async Task<Message> GetAsync(object id, bool noTracking = true)
|
|
{
|
|
return await Repository.SingleOrDefaultAsync(c => c.Id == id.ToString(), noTracking).ConfigureAwait(false);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Liste der Nachrichten für eine Konversation zurück
|
|
/// </summary>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <param name="senderId">Id des Senders</param>
|
|
/// <returns>Liste Nachrichten</returns>
|
|
public async Task<List<Message>> 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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Liste der Nachrichten für eine Konversation zurück
|
|
/// </summary>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <param name="take">Wie viele News sollen abgerufen werden? -1 Wenn nicht anwenden.</param>
|
|
/// <param name="skip">Wie viele News sollen ausgelassen werden? -1 Wenn nicht anwenden.</param>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <param name="senderId">Id des Senders</param>
|
|
/// <returns>Liste Nachrichten</returns>
|
|
public async Task<ListCommunicationResult<List<Message>>> GetMessagesAsync(string senderId, string conversationId, int take, int skip, string accessToken, CancellationToken token)
|
|
{
|
|
var result = new ListCommunicationResult<List<Message>>() { Value = new List<Message>() };
|
|
|
|
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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Liste der Nachrichten für die Synchronisation zurück
|
|
/// </summary>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <param name="senderId">Id des Senders</param>
|
|
/// <returns>Liste Nachrichten</returns>
|
|
public async Task<CommunicationResult<List<Message>>> GetMessagesForSyncAsync(string senderId, string accessToken, CancellationToken token)
|
|
{
|
|
var result = new CommunicationResult<List<Message>>() { Value = new List<Message>() };
|
|
|
|
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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Erstellen einer Nachricht. Wird lokal gespeichert und an der Server übertragen.
|
|
/// Wenn keine Verbindung besteht, wird ein Delta erstellt
|
|
/// </summary>
|
|
/// <param name="senderId">id des Senders</param>
|
|
/// <param name="receiverId">Id des Empfängers</param>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <param name="content">Inhalt</param>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <returns>Task</returns>
|
|
public async Task<Message> 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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Setzt alle Nachrichten in einer Konversation als gelesen
|
|
/// </summary>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <returns>Task</returns>
|
|
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();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Liste ungelesener Nachrichten zurück
|
|
/// </summary>
|
|
/// <returns>Liste ungelesener Nachrichten</returns>
|
|
public async Task<List<Message>> GetUnreadAsync()
|
|
{
|
|
var unreadMessages = await Repository.FindAsync(c => c.Read == false).ConfigureAwait(false);
|
|
return unreadMessages.ToList();
|
|
}
|
|
|
|
/// <summary>
|
|
/// Gibt eine Nachricht einer Konversation zurück
|
|
/// </summary>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <param name="messageId">Id der Nachricht</param>
|
|
/// <returns>Nachricht oder null, wenn nicht gefunden</returns>
|
|
public async Task<Message> GetAsync(string conversationId, string messageId)
|
|
{
|
|
return await Repository.FirstOrDefaultAsync(c => c.ConversationId == conversationId && c.Id == messageId).ConfigureAwait(false);
|
|
}
|
|
|
|
/// <summary>
|
|
/// Löscht eine Nachricht.
|
|
/// Die Nachricht wird auf gelöscht gesetzt.
|
|
/// </summary>
|
|
/// <param name="conversationId">Id der Konversation</param>
|
|
/// <param name="messageId">Id der Nachricht</param>
|
|
/// <param name="forceDelete">True wenn wirklich gelöscht werden soll</param>
|
|
/// <returns>True wenn erfolgreich, false sonst</returns>
|
|
public async Task<bool> 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
|
|
|
|
/// <summary>
|
|
/// Holen der Daten vom lokalen Speicher die noch nicht synchronisiert wurden und senden an den Server
|
|
/// </summary>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <returns>Task</returns>
|
|
async Task<SyncResult<Message>> ISyncPushPullService<Message>.PushAsync(string accessToken, CancellationToken token)
|
|
{
|
|
var result = new SyncResult<Message>();
|
|
|
|
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<Message>(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;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Holen der letzten Daten vom Server und Synchronisieren mit den lokalen Daten
|
|
/// </summary>
|
|
/// <param name="language">Sprache</param>
|
|
/// <param name="accessToken">Aktuelles Accesstoken</param>
|
|
/// <param name="token">CancellationToken</param>
|
|
/// <returns>Task</returns>
|
|
async Task<SyncResult<Message>> ISyncPushPullService<Message>.PullAsync(string language, string accessToken, CancellationToken token)
|
|
{
|
|
var result = new SyncResult<Message>();
|
|
|
|
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
|
|
|
|
/// <summary>
|
|
/// Behandeln der Liste von Nachrichten wenn welche vom Online-Store geholt werden.
|
|
/// </summary>
|
|
/// <param name="messages">Liste der Konversationen</param>
|
|
/// <returns>Task</returns>
|
|
private async Task<SyncResult<Message>> HandleChangesAsync(List<Message> messages)
|
|
{
|
|
var result = new SyncResult<Message>();
|
|
|
|
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
|
|
}
|
|
}
|