using MassTransit; using OpenWarehouse.warehouse.api.Common.ItemTagManager; using OpenWarehouse.warehouse.api.Common.Kafka; using OpenWarehouse.warehouse.api.Common.Kafka.TagTree; using OpenWarehouse.warehouse.api.Common.TaskManager; using OpenWarehouse.warehouse.api.Model.Item; using ItemCategoryTag = OpenWarehouse.warehouse.api.Model.Item.ItemCategoryTag; using TaskStatus = OpenWarehouse.warehouse.api.Common.TaskManager.TaskStatus; namespace OpenWarehouse.warehouse.api.Services; public class TagTreeConsumer(ILogger logger, ITaskManager taskManager ,IItemTagManager itemTagManager) : IConsumer { public async Task Consume(ConsumeContext context) { logger.LogInformation("Received and processing Kafka Message: {Guid}", context.Message.Guid); var status = TaskStatus.Processing; try { await taskManager.UpdateTaskStatus(context.Message.Guid, status); switch (context.Message.Action) { case Actions.Create: status = await Create(context.Message); break; case Actions.Delete: status = await Delete(context.Message.ItemCategoryTagsList); break; default: logger.LogWarning("Action not allowed in request: {Guid}", context.Message.Guid); status = TaskStatus.Error; break; } await taskManager.UpdateTaskStatus(context.Message.Guid, status); } catch (Exception ex) { logger.LogError(ex, "Error occurred while processing Kafka message: {Guid}", context.Message.Guid); status = TaskStatus.Error; await taskManager.UpdateTaskStatus(context.Message.Guid, status); } } private async Task Create(TagTreeKafkaMessage context) { try { var len = context.ItemCategoryTagsList.Count; var preProcessingData = context.ItemCategoryTagsList .Distinct() .ToList(); var result = await itemTagManager.CreateTagTree(preProcessingData); if (result == 0) { logger.LogWarning("Failed to create tag tree for message: {Guid}", context.Guid); return TaskStatus.Error; } else if (result < len) { logger.LogInformation("Partial success in creating tag tree for message: {Guid}", context.Guid); return TaskStatus.NotFullSuccess; } else { logger.LogInformation("Successfully created tag tree for message: {Guid}", context.Guid); return TaskStatus.Success; } } catch (Exception ex) { logger.LogError(ex, "Error occurred while creating tag tree for message: {Guid}", context.Guid); return TaskStatus.Error; } } private async Task Delete(List names) { try { var result = await itemTagManager.DeleteTagsByNames(names); if (result == 0) { logger.LogWarning("Failed to delete tags for names: {Names}", names); return TaskStatus.Error; } else if (result < names.Count) { logger.LogInformation("Partial success in deleting tags for names: {Names}", names); return TaskStatus.NotFullSuccess; } else { logger.LogInformation("Successfully deleted tags for names: {Names}", names); return TaskStatus.Success; } } catch (Exception ex) { logger.LogError(ex, "Error occurred while deleting tags for names: {Names}", names); return TaskStatus.Error; } } }